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

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

Returns:

  • (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)
  transport_message_class(: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.message_context,
    trace_context: envelope.trace_context,
    code: error_code(error),
    message: error.message,
    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)
  transport_message_class(: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.message_context,
    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

Returns:

  • (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.

Returns:

  • (Boolean)


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 transport_message_class(name)
  ::Legion::Extensions::Llm::Transport::Messages.const_get(name)
rescue LoadError, NameError => e
  raise ConfigurationError, "fleet responder transport unavailable for #{name}: #{e.message}"
end

.truthy?(value) ⇒ Boolean

Returns:

  • (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.

Raises:



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.

Raises:

  • (ArgumentError)


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

Raises:

  • (ArgumentError)


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