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. Protocol v3 (06 W9): parse → legacy-field rejection → required fields (one Protocol::REQUIRED_FIELDS list) → explicit version → provider-family match → exact execution contract (required, P2) → WorkerExecution.call → publish FleetResponse (E3/E4) → ack. On error: publish FleetError (E6), reject per F6, re-raise.
Defined Under Namespace
Classes: ConfigurationError
Class Method Summary collapse
- .ack(delivery) ⇒ Object
-
.call(payload:, provider_family:, 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
- .enabled_for?(provider_instances) ⇒ Boolean
- .error_code(error) ⇒ Object
-
.parse_json(payload) ⇒ Object
L1: Legion::JSON only (house rule) — the bare ::JSON fallback is deleted; Legion::JSON is a hard dependency of this gem.
-
.parse_payload(payload) ⇒ Object
E1: the single wire-normalization entry.
-
.publish_error(envelope, error) ⇒ Object
E6: the error envelope.
-
.publish_response(envelope, response) ⇒ Object
E3: the response envelope carries the serialized Canonical::Response.
- .reject(delivery, requeue:) ⇒ Object
- .reject_legacy_fields!(envelope) ⇒ Object
- .requeue_error?(error) ⇒ Boolean
- .resolve_provider_instances(provider_instances) ⇒ Object
-
.retryable_error?(error) ⇒ Boolean
F6: retryability is derived from the one ProviderOutcome kind table (05 O6) — transient kinds retry; contract/policy/auth/classification kinds never do.
- .safe_publish_error(envelope, error) ⇒ Object
- .transport_message_class(name) ⇒ Object
- .truthy?(value) ⇒ Boolean
-
.validate_execution_contract!(envelope) ⇒ Object
P2: exact execution only — the marker is required and must equal the exact marker; absence is rejected.
-
.validate_protocol_version!(envelope) ⇒ Object
P3: explicit version — no default fill.
- .validate_provider_family!(envelope, provider_family) ⇒ Object
Class Method Details
.ack(delivery) ⇒ Object
165 166 167 168 169 170 171 172 173 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 165 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 |
.call(payload:, provider_family:, 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. L6: the dead provider_class/provider_instances params are deleted — v3 dispatch is exact-only and never constructs a provider; passing provider objects here was a latent second execution truth.
45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 45 def call(payload:, provider_family:, registry: ::Legion::Extensions::Llm::Inventory::Registry, delivery: nil, properties: nil) envelope = parse_payload(payload) check_envelope!(envelope, provider_family:) response = WorkerExecution.call(envelope: envelope, 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
81 82 83 84 85 86 87 88 89 90 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 81 def check_envelope!(envelope, provider_family:) reject_legacy_fields!(envelope) Protocol::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 |
.enabled_for?(provider_instances) ⇒ Boolean
61 62 63 64 65 66 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 61 def enabled_for?(provider_instances) instances = resolve_provider_instances(provider_instances) instances.any? do |_instance_id, settings| truthy?(Utils.deep_symbolize_keys(settings).dig(:fleet, :respond_to_requests)) end end |
.error_code(error) ⇒ Object
235 236 237 238 239 240 241 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 235 def error_code(error) return 'configuration_error' if error.is_a?(ConfigurationError) return 'policy_error' if error.is_a?(WorkerExecution::PolicyError) return 'contract_error' if error.is_a?(ContractError) 'provider_error' end |
.parse_json(payload) ⇒ Object
L1: Legion::JSON only (house rule) — the bare ::JSON fallback is deleted; Legion::JSON is a hard dependency of this gem.
187 188 189 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 187 def parse_json(payload) ::Legion::JSON.parse(payload) end |
.parse_payload(payload) ⇒ Object
E1: the single wire-normalization entry. Wrong-shape payloads (neither Hash, String, nor envelope) raise — the silent {} fallback is deleted.
71 72 73 74 75 76 77 78 79 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 71 def parse_payload(payload) case payload when FleetEnvelope then payload when String then FleetEnvelope.new(data: parse_json(payload)) when Hash then FleetEnvelope.new(data: payload) else raise ContractError, "fleet payload expected Hash or String, got #{payload.class}" end end |
.publish_error(envelope, error) ⇒ Object
E6: the error envelope. retryable is derived per F6.
127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 127 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: envelope.execution_contract, offering_id: envelope.offering_id ).publish end |
.publish_response(envelope, response) ⇒ Object
E3: the response envelope carries the serialized Canonical::Response. E4/G5: thinking never crosses the fleet — excluded exactly once, here.
107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 107 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, response: response.to_h.except(:thinking), execution_contract: envelope.execution_contract, offering_id: envelope.offering_id ).publish end |
.reject(delivery, requeue:) ⇒ Object
175 176 177 178 179 180 181 182 183 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 175 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
191 192 193 194 195 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 191 def reject_legacy_fields!(envelope) Protocol::LEGACY_FIELDS.each do |field| raise ArgumentError, "#{field} is not supported by fleet protocol v3" if envelope.key?(field) end end |
.requeue_error?(error) ⇒ Boolean
215 216 217 218 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 215 def requeue_error?(error) retryable_error?(error) && Settings.value(:fleet, :consumer, :requeue_transient, default: true) != false end |
.resolve_provider_instances(provider_instances) ⇒ Object
210 211 212 213 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 210 def resolve_provider_instances(provider_instances) instances = provider_instances.respond_to?(:call) ? provider_instances.call : provider_instances Utils.deep_symbolize_keys(instances || {}) end |
.retryable_error?(error) ⇒ Boolean
F6: retryability is derived from the one ProviderOutcome kind table (05 O6) — transient kinds retry; contract/policy/auth/classification kinds never do. The default-true catch-all is deleted.
223 224 225 226 227 228 229 230 231 232 233 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 223 def retryable_error?(error) return false if error.is_a?(ConfigurationError) return false if error.is_a?(WorkerExecution::PolicyError) return false if error.is_a?(ContractError) return false if error.is_a?(TokenError) return false if error.is_a?(Inventory::Errors::ExactOfferingMismatchError) ::Legion::Extensions::Llm::Routing::ProviderOutcome::RETRYABLE_KINDS.include?( ::Legion::Extensions::Llm::Routing::ProviderOutcome.kind_for(error) ) end |
.safe_publish_error(envelope, error) ⇒ Object
150 151 152 153 154 155 156 157 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 150 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
159 160 161 162 163 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 159 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
243 244 245 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 243 def truthy?(value) value == true || value.to_s == 'true' end |
.validate_execution_contract!(envelope) ⇒ Object
P2: exact execution only — the marker is required and must equal the exact marker; absence is rejected. The exact fields are additionally required by the marker.
95 96 97 98 99 100 101 102 103 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 95 def validate_execution_contract!(envelope) marker = envelope.execution_contract raise ContractError, "execution_contract must be #{Protocol::EXACT_EXECUTION_CONTRACT}" unless marker == Protocol::EXACT_EXECUTION_CONTRACT Protocol::EXACT_REQUIRED_FIELDS.each do |field| raise ContractError, "#{field} is required for #{Protocol::EXACT_EXECUTION_CONTRACT}" unless envelope.key?(field) && !envelope[field].nil? end end |
.validate_protocol_version!(envelope) ⇒ Object
P3: explicit version — no default fill.
198 199 200 201 202 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 198 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
204 205 206 207 208 |
# File 'lib/legion/extensions/llm/fleet/provider_responder.rb', line 204 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 |