Module: Legion::Extensions::Llm::Fleet::ProviderResponder
- Extended by:
- Logging::Helper
- Includes:
- Logging::Helper
- Defined in:
- lib/legion/extensions/llm/fleet/provider_responder.rb
Overview
Shared implementation for provider-owned fleet responder runners.
Defined Under Namespace
Classes: ConfigurationError, FleetEnvelope
Constant Summary collapse
- REQUIRED_FIELDS =
%i[ request_id correlation_id idempotency_key operation provider provider_instance model params reply_to message_context caller trace_context signed_token timeout_seconds expires_at protocol_version ].freeze
- LEGACY_FIELDS =
%i[schema_version request_type fleet_correlation_id].freeze
Class Method Summary collapse
- .ack(delivery) ⇒ Object
- .build_provider(envelope:, provider_class:, provider_instances:) ⇒ Object
-
.call(payload:, provider_family:, provider_class:, provider_instances:, registry: ::Legion::Extensions::Llm::Inventory::Registry, delivery: nil, properties: nil) ⇒ Object
Public runner entry point mirrors AMQP delivery callbacks, which carry both delivery and property metadata.
- .check_envelope!(envelope, provider_family:) ⇒ Object
- .deep_symbolize(value) ⇒ Object
- .dig(hash, *keys) ⇒ Object
-
.dispatch_request(envelope, provider_class, provider_instances, registry) ⇒ Object
Exact requests dispatch through the registry and never call build_provider; legacy v2 keeps the provider-object path.
-
.enabled_for?(provider_instances) ⇒ Boolean
rubocop:enable Metrics/ParameterLists.
- .error_code(error) ⇒ Object
- .exact?(envelope) ⇒ Boolean
- .parse_json(payload) ⇒ Object
- .parse_payload(payload) ⇒ Object
- .publish_error(envelope, error) ⇒ Object
- .publish_response(envelope, response) ⇒ Object
- .reject(delivery, requeue:) ⇒ Object
- .reject_legacy_fields!(envelope) ⇒ Object
- .requeue_error?(error) ⇒ Boolean
- .resolve_provider_instances(provider_instances) ⇒ Object
- .response_content(response) ⇒ Object
- .response_field(response, field) ⇒ Object
- .response_metadata(response) ⇒ Object
- .response_usage(response) ⇒ Object
- .retryable_error?(error) ⇒ Boolean
- .safe_publish_error(envelope, error) ⇒ Object
- .transport_message_class(name) ⇒ Object
- .truthy?(value) ⇒ Boolean
-
.validate_execution_contract!(envelope) ⇒ Object
Marker absence means legacy v2; an unknown nonempty marker is rejected; the exact marker additionally requires every EXACT_REQUIRED_FIELDS value.
- .validate_protocol_version!(envelope) ⇒ Object
- .validate_provider_family!(envelope, provider_family) ⇒ Object
Class Method Details
.ack(delivery) ⇒ Object
219 220 221 222 223 224 225 226 227 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 219 def ack(delivery) return unless delivery if delivery.respond_to?(:ack) delivery.ack elsif delivery.respond_to?(:channel) && delivery.respond_to?(:delivery_tag) delivery.channel.ack(delivery.delivery_tag) end end |
.build_provider(envelope:, provider_class:, provider_instances:) ⇒ Object
145 146 147 148 149 150 151 152 153 154 155 156 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 145 def build_provider(envelope:, provider_class:, provider_instances:) instances = resolve_provider_instances(provider_instances) instance_id = envelope.provider_instance.to_s instance_settings = instances[instance_id.to_sym] || instances[instance_id] unless instance_settings raise ConfigurationError, "fleet provider instance is not configured: #{instance_id}" end raise ConfigurationError, "fleet responses are disabled for provider instance: #{instance_id}" unless truthy?(dig(instance_settings, :fleet, :respond_to_requests)) provider_class.new(deep_symbolize(instance_settings)) end |
.call(payload:, provider_family:, provider_class:, provider_instances:, registry: ::Legion::Extensions::Llm::Inventory::Registry, delivery: nil, properties: nil) ⇒ Object
Public runner entry point mirrors AMQP delivery callbacks, which carry both delivery and property metadata. rubocop:disable Metrics/ParameterLists
71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 71 def call(payload:, provider_family:, provider_class:, provider_instances:, registry: ::Legion::Extensions::Llm::Inventory::Registry, delivery: nil, properties: nil) envelope = parse_payload(payload) check_envelope!(envelope, provider_family:) response = dispatch_request(envelope, provider_class, provider_instances, registry) publish_response(envelope, response) ack(delivery || properties) response rescue StandardError => e handle_exception(e, level: :warn, handled: false, operation: 'llm.fleet.provider_responder.call', provider_family:) safe_publish_error(envelope, e) if defined?(envelope) && envelope reject(delivery || properties, requeue: requeue_error?(e)) raise end |
.check_envelope!(envelope, provider_family:) ⇒ Object
107 108 109 110 111 112 113 114 115 116 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 107 def check_envelope!(envelope, provider_family:) reject_legacy_fields!(envelope) REQUIRED_FIELDS.each do |field| raise ArgumentError, "#{field} is required" unless envelope.key?(field) && !envelope[field].nil? end validate_protocol_version!(envelope) validate_provider_family!(envelope, provider_family) validate_execution_contract!(envelope) end |
.deep_symbolize(value) ⇒ Object
329 330 331 332 333 334 335 336 337 338 339 340 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 329 def deep_symbolize(value) case value when Hash value.each_with_object({}) do |(key, child), result| result[key.respond_to?(:to_sym) ? key.to_sym : key] = deep_symbolize(child) end when Array value.map { |child| deep_symbolize(child) } else value end end |
.dig(hash, *keys) ⇒ Object
317 318 319 320 321 322 323 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 317 def dig(hash, *keys) keys.reduce(hash) do |current, key| break nil unless current.respond_to?(:key?) current[key.to_sym] || current[key.to_s] end end |
.dispatch_request(envelope, provider_class, provider_instances, registry) ⇒ Object
Exact requests dispatch through the registry and never call build_provider; legacy v2 keeps the provider-object path.
136 137 138 139 140 141 142 143 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 136 def dispatch_request(envelope, provider_class, provider_instances, registry) if exact?(envelope) WorkerExecution.call(envelope: envelope, registry: registry) else provider = build_provider(envelope:, provider_class:, provider_instances:) WorkerExecution.call(envelope: envelope, provider: provider) end end |
.enabled_for?(provider_instances) ⇒ Boolean
rubocop:enable Metrics/ParameterLists
88 89 90 91 92 93 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 88 def enabled_for?(provider_instances) instances = resolve_provider_instances(provider_instances) instances.any? do |_instance_id, settings| truthy?(dig(settings, :fleet, :respond_to_requests)) end end |
.error_code(error) ⇒ Object
282 283 284 285 286 287 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 282 def error_code(error) return 'configuration_error' if error.is_a?(ConfigurationError) return 'policy_error' if error.is_a?(WorkerExecution::PolicyError) 'provider_error' end |
.exact?(envelope) ⇒ Boolean
130 131 132 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 130 def exact?(envelope) envelope.execution_contract == Protocol::EXACT_EXECUTION_CONTRACT end |
.parse_json(payload) ⇒ Object
239 240 241 242 243 244 245 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 239 def parse_json(payload) if defined?(::Legion::JSON) ::Legion::JSON.parse(payload) else ::JSON.parse(payload) end end |
.parse_payload(payload) ⇒ Object
95 96 97 98 99 100 101 102 103 104 105 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 95 def parse_payload(payload) hash = case payload when FleetEnvelope payload.to_h when String parse_json(payload) else payload.respond_to?(:to_h) ? payload.to_h : {} end FleetEnvelope.new(data: deep_symbolize(hash)) end |
.publish_error(envelope, error) ⇒ Object
181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 181 def publish_error(envelope, error) (:FleetError).new( protocol_version: envelope.protocol_version, request_id: envelope.request_id, correlation_id: envelope.correlation_id, idempotency_key: envelope.idempotency_key, operation: envelope.operation, provider: envelope.provider, provider_instance: envelope.provider_instance, model: envelope.model, reply_to: envelope.reply_to, message_context: envelope., trace_context: envelope.trace_context, code: error_code(error), message: error., error_class: error.class.name, retryable: retryable_error?(error), metadata: {}, execution_contract: exact?(envelope) ? envelope.execution_contract : nil, offering_id: exact?(envelope) ? envelope.offering_id : nil ).publish end |
.publish_response(envelope, response) ⇒ Object
158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 158 def publish_response(envelope, response) (:FleetResponse).new( protocol_version: envelope.protocol_version, request_id: envelope.request_id, correlation_id: envelope.correlation_id, idempotency_key: envelope.idempotency_key, operation: envelope.operation, provider: envelope.provider, provider_instance: envelope.provider_instance, model: envelope.model, reply_to: envelope.reply_to, message_context: envelope., trace_context: envelope.trace_context, content: response_content(response), tool_calls: response_field(response, :tool_calls) || [], usage: response_usage(response), finish_reason: response_field(response, :finish_reason), metadata: (response), execution_contract: exact?(envelope) ? envelope.execution_contract : nil, offering_id: exact?(envelope) ? envelope.offering_id : nil ).publish end |
.reject(delivery, requeue:) ⇒ Object
229 230 231 232 233 234 235 236 237 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 229 def reject(delivery, requeue:) return unless delivery if delivery.respond_to?(:reject) delivery.reject(requeue) elsif delivery.respond_to?(:channel) && delivery.respond_to?(:delivery_tag) delivery.channel.reject(delivery.delivery_tag, requeue) end end |
.reject_legacy_fields!(envelope) ⇒ Object
247 248 249 250 251 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 247 def reject_legacy_fields!(envelope) LEGACY_FIELDS.each do |field| raise ArgumentError, "#{field} is not supported by fleet protocol v2" if envelope.key?(field) end end |
.requeue_error?(error) ⇒ Boolean
270 271 272 273 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 270 def requeue_error?(error) retryable_error?(error) && Settings.value(:fleet, :consumer, :requeue_transient, default: true) != false end |
.resolve_provider_instances(provider_instances) ⇒ Object
265 266 267 268 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 265 def resolve_provider_instances(provider_instances) instances = provider_instances.respond_to?(:call) ? provider_instances.call : provider_instances deep_symbolize(instances || {}) end |
.response_content(response) ⇒ Object
289 290 291 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 289 def response_content(response) response_field(response, :content) || response_field(response, :result) || response.to_s end |
.response_field(response, field) ⇒ Object
309 310 311 312 313 314 315 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 309 def response_field(response, field) return response[field] if response.respond_to?(:key?) && response.key?(field) return response[field.to_s] if response.respond_to?(:key?) && response.key?(field.to_s) return response.public_send(field) if response.respond_to?(field) nil end |
.response_metadata(response) ⇒ Object
304 305 306 307 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 304 def (response) = response_field(response, :metadata) .respond_to?(:to_h) ? deep_symbolize() : {} end |
.response_usage(response) ⇒ Object
293 294 295 296 297 298 299 300 301 302 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 293 def response_usage(response) usage = response_field(response, :usage) || response_field(response, :tokens) return deep_symbolize(usage) if usage.respond_to?(:to_h) { input_tokens: response_field(response, :input_tokens), output_tokens: response_field(response, :output_tokens), thinking_tokens: response_field(response, :thinking_tokens) }.compact end |
.retryable_error?(error) ⇒ Boolean
275 276 277 278 279 280 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 275 def retryable_error?(error) return false if error.is_a?(ConfigurationError) return false if error.is_a?(WorkerExecution::PolicyError) true end |
.safe_publish_error(envelope, error) ⇒ Object
204 205 206 207 208 209 210 211 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 204 def safe_publish_error(envelope, error) publish_error(envelope, error) rescue StandardError => e handle_exception(e, level: :warn, handled: true, operation: 'llm.fleet.provider_responder.safe_publish_error', error_class: error.class.name) nil end |
.transport_message_class(name) ⇒ Object
213 214 215 216 217 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 213 def (name) ::Legion::Extensions::Llm::Transport::Messages.const_get(name) rescue LoadError, NameError => e raise ConfigurationError, "fleet responder transport unavailable for #{name}: #{e.}" end |
.truthy?(value) ⇒ Boolean
325 326 327 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 325 def truthy?(value) value == true || value.to_s == 'true' end |
.validate_execution_contract!(envelope) ⇒ Object
Marker absence means legacy v2; an unknown nonempty marker is rejected; the exact marker additionally requires every EXACT_REQUIRED_FIELDS value.
120 121 122 123 124 125 126 127 128 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 120 def validate_execution_contract!(envelope) marker = envelope.execution_contract return if marker.nil? raise ArgumentError, "unknown execution_contract: #{marker}" unless marker == Protocol::EXACT_EXECUTION_CONTRACT Protocol::EXACT_REQUIRED_FIELDS.each do |field| raise ArgumentError, "#{field} is required for #{Protocol::EXACT_EXECUTION_CONTRACT}" unless envelope.key?(field) && !envelope[field].nil? end end |
.validate_protocol_version!(envelope) ⇒ Object
253 254 255 256 257 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 253 def validate_protocol_version!(envelope) return if envelope.protocol_version == Protocol::VERSION raise ArgumentError, "protocol_version must be #{Protocol::VERSION}" end |
.validate_provider_family!(envelope, provider_family) ⇒ Object
259 260 261 262 263 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 259 def validate_provider_family!(envelope, provider_family) return if envelope.provider.to_s == provider_family.to_s raise ArgumentError, "fleet request provider #{envelope.provider} does not match #{provider_family}" end |