Module: Legion::LLM::Metering

Extended by:
Legion::Logging::Helper
Defined in:
lib/legion/llm/metering.rb,
lib/legion/llm/metering/tokens.rb,
lib/legion/llm/metering/tracker.rb,
lib/legion/llm/metering/estimator.rb

Defined Under Namespace

Modules: Pricing, Recorder, Tokens

Constant Summary collapse

SPOOL_DIR =
File.expand_path('~/.legionio/data/spool/metering')
SPOOL_FILE =
File.join(SPOOL_DIR, 'events.jsonl').freeze
SPOOL_MUTEX =
Mutex.new

Class Method Summary collapse

Class Method Details

.attributed_event(event) ⇒ Object



47
48
49
50
51
52
# File 'lib/legion/llm/metering.rb', line 47

def attributed_event(event)
  source = event.is_a?(Hash) ? event.dup : {}
  source[:identity] ||= Legion::LLM::PublisherIdentity.current
  source[:caller] ||= Legion::LLM::PublisherIdentity.caller_hash
  source
end

.const_missing(name) ⇒ Object

Backward-compat: resolve old Legion::LLM::Metering::Exchange, ::Event



298
299
300
301
302
303
304
305
306
307
308
309
# File 'lib/legion/llm/metering.rb', line 298

def self.const_missing(name)
  case name
  when :Exchange
    require_relative 'transport/exchanges/metering'
    Transport::Exchanges::Metering
  when :Event
    require_relative 'transport/messages/metering_event'
    Transport::Messages::MeteringEvent
  else
    super
  end
end

.crypt_available?Boolean

Returns:

  • (Boolean)


219
220
221
# File 'lib/legion/llm/metering.rb', line 219

def crypt_available?
  defined?(Legion::Crypt) && Legion::Crypt.respond_to?(:encrypt)
end

.decrypt_spool_line(line) ⇒ Object



223
224
225
226
227
228
229
230
231
232
233
234
# File 'lib/legion/llm/metering.rb', line 223

def decrypt_spool_line(line)
  return line unless defined?(Legion::Crypt) && Legion::Crypt.respond_to?(:decrypt)
  return line if line.start_with?('{')

  Legion::Crypt.decrypt(line)
rescue StandardError => e
  # N5: a corrupt line must not wedge the flush silently — log the
  # fault and drop just this line (nil) so the rest of the batch
  # still drains.
  handle_exception(e, level: :error, operation: 'llm.metering.decrypt_spool_line')
  nil
end

.emit(event) ⇒ Object



29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
# File 'lib/legion/llm/metering.rb', line 29

def emit(event)
  event = attributed_event(event)
  event_class = metering_event_class if transport_connected?

  if event_class
    event_class.new(**event).publish
    log.info("[llm][metering] published provider=#{event[:provider]} model=#{event[:model_id]}")
    :published
  else
    spool_event(event)
    log.info("[llm][metering] spooled provider=#{event[:provider]} model=#{event[:model_id]} reason=transport_unavailable")
    :spooled
  end
rescue StandardError => e
  handle_exception(e, level: :warn, operation: 'llm.metering.emit')
  :dropped
end

.encrypt_spool?Boolean

Returns:

  • (Boolean)


208
209
210
211
212
213
214
215
216
217
# File 'lib/legion/llm/metering.rb', line 208

def encrypt_spool?
  return false unless crypt_available?

  Legion::Settings.dig(:llm, :compliance, :encrypt_spool) == true
rescue StandardError => e
  # No-fail-open: a settings fault must never silently disable spool
  # encryption. Fail closed: encrypt whenever Crypt is available.
  handle_exception(e, level: :error, operation: 'llm.metering.encrypt_spool?')
  crypt_available?
end

.enforce_max_eventsObject



262
263
264
265
266
267
268
269
270
271
272
273
274
# File 'lib/legion/llm/metering.rb', line 262

def enforce_max_events
  path = spool_file_path
  return unless File.exist?(path)

  max = spool_settings[:max_events] || 10_000
  lines = File.readlines(path, chomp: true)
  return if lines.size < max

  # Drop oldest events to make room
  trimmed = lines.last(max - 1)
  File.write(path, trimmed.map { |l| "#{l}\n" }.join)
  log.debug("[llm][metering] enforce_max_events trimmed=#{lines.size - trimmed.size} max=#{max}")
end

.ensure_spool_dirObject



276
277
278
# File 'lib/legion/llm/metering.rb', line 276

def ensure_spool_dir
  FileUtils.mkdir_p(spool_dir_path)
end

.extract_dispatch_path(response) ⇒ Object



349
350
351
352
353
# File 'lib/legion/llm/metering.rb', line 349

def extract_dispatch_path(response)
  return nil unless response.is_a?(Hash)

  response[:dispatch_path] || response[:tier] || response.dig(:routing, :tier)
end

.extract_error_category(response) ⇒ Object



325
326
327
328
329
# File 'lib/legion/llm/metering.rb', line 325

def extract_error_category(response)
  return nil unless response.is_a?(Hash)

  response[:error_category] || response.dig(:error, :category)
end

.extract_error_code(response) ⇒ Object



331
332
333
334
335
# File 'lib/legion/llm/metering.rb', line 331

def extract_error_code(response)
  return nil unless response.is_a?(Hash)

  response[:error_code] || response.dig(:error, :code)
end

.extract_error_message(response) ⇒ Object



337
338
339
340
341
# File 'lib/legion/llm/metering.rb', line 337

def extract_error_message(response)
  return nil unless response.is_a?(Hash)

  response[:error_message] || response.dig(:error, :message)
end

.extract_finish_reason(response) ⇒ Object

--- Extractor helpers for after_chat hook ---



313
314
315
316
317
# File 'lib/legion/llm/metering.rb', line 313

def extract_finish_reason(response)
  return nil unless response.is_a?(Hash)

  response[:finish_reason] || response.dig(:stop, :reason) || response.dig(:choices, 0, :finish_reason)
end

.extract_hash_value(hash, key) ⇒ Object



180
181
182
183
184
185
186
187
# File 'lib/legion/llm/metering.rb', line 180

def extract_hash_value(hash, key)
  return nil unless hash.respond_to?(:key?)

  string_key = key.to_s
  return hash[string_key] if hash.key?(string_key)

  hash[key] if hash.key?(key)
end

.extract_latency_ms(response) ⇒ Object



361
362
363
364
365
# File 'lib/legion/llm/metering.rb', line 361

def extract_latency_ms(response)
  return nil unless response.is_a?(Hash)

  (response[:latency_ms] || response.dig(:timing, :latency_ms) || 0).to_i
end

.extract_model(response) ⇒ Object



174
175
176
177
178
# File 'lib/legion/llm/metering.rb', line 174

def extract_model(response)
  return nil unless response.is_a?(Hash)

  extract_hash_value(extract_hash_value(response, :meta), :model) || extract_hash_value(response, :model)
end

.extract_provider(response) ⇒ Object



168
169
170
171
172
# File 'lib/legion/llm/metering.rb', line 168

def extract_provider(response)
  return nil unless response.is_a?(Hash)

  extract_hash_value(extract_hash_value(response, :meta), :provider) || extract_hash_value(response, :provider)
end

.extract_provider_instance(response) ⇒ Object



343
344
345
346
347
# File 'lib/legion/llm/metering.rb', line 343

def extract_provider_instance(response)
  return nil unless response.is_a?(Hash)

  response[:provider_instance] || response.dig(:routing, :provider_instance) || response[:instance]
end

.extract_route_attempts(response) ⇒ Object



355
356
357
358
359
# File 'lib/legion/llm/metering.rb', line 355

def extract_route_attempts(response)
  return nil unless response.is_a?(Hash)

  (response[:route_attempts] || 0).to_i
end

.extract_status(response) ⇒ Object



373
374
375
376
377
# File 'lib/legion/llm/metering.rb', line 373

def extract_status(response)
  return 'success' unless response.is_a?(Hash)

  response[:error] ? 'failure' : 'success'
end

.extract_tier(response) ⇒ Object



319
320
321
322
323
# File 'lib/legion/llm/metering.rb', line 319

def extract_tier(response)
  return nil unless response.is_a?(Hash)

  response[:tier] || response.dig(:routing, :tier)
end

.extract_usage(response) ⇒ Object



158
159
160
161
162
163
164
165
166
# File 'lib/legion/llm/metering.rb', line 158

def extract_usage(response)
  return { input_tokens: 0, output_tokens: 0 } unless response.is_a?(Hash)

  usage = extract_hash_value(response, :usage) || {}
  {
    input_tokens:  extract_hash_value(usage, :input_tokens) || extract_hash_value(usage, :prompt_tokens) || 0,
    output_tokens: extract_hash_value(usage, :output_tokens) || extract_hash_value(usage, :completion_tokens) || 0
  }
end

.extract_wall_clock_ms(response) ⇒ Object



367
368
369
370
371
# File 'lib/legion/llm/metering.rb', line 367

def extract_wall_clock_ms(response)
  return nil unless response.is_a?(Hash)

  (response[:wall_clock_ms] || response.dig(:timing, :wall_clock_ms) || 0).to_i
end

.flush_spoolObject



54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
# File 'lib/legion/llm/metering.rb', line 54

def flush_spool
  return 0 unless File.exist?(spool_file_path)

  event_class = metering_event_class
  unless event_class && transport_connected?
    log.debug('[llm][metering] flush_spool skipped reason=transport_unavailable')
    return 0
  end

  # Read and truncate atomically under the mutex so no events written
  # between read and truncate can be silently lost.
  events = SPOOL_MUTEX.synchronize do
    path = spool_file_path
    return 0 unless File.exist?(path)

    lines = File.readlines(path, chomp: true)
    parsed = lines.filter_map do |line|
      next if line.strip.empty?

      decrypted = decrypt_spool_line(line)
      next if decrypted.nil?

      Legion::JSON.load(decrypted)
    end
    File.write(path, '')
    parsed
  end

  return 0 if events.empty?

  batch_sleep = spool_settings[:flush_batch_sleep] || 0.0
  flushed = 0

  events.each_with_index do |event_data, index|
    event_class.new(**event_data).publish
    flushed += 1
    sleep(batch_sleep) if batch_sleep.positive? && index < events.size - 1
  end

  log.info("[llm][metering] flush_spool flushed=#{flushed}")
  flushed
rescue StandardError => e
  handle_exception(e, level: :warn, operation: 'llm.metering.flush_spool')
  0
end

.install_hookObject



100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
# File 'lib/legion/llm/metering.rb', line 100

def install_hook
  Legion::LLM::Hooks.after_chat do |response:, model:, caller: nil, **|
    usage = extract_usage(response)
    next if usage[:input_tokens].zero? && usage[:output_tokens].zero?

    resolved_model    = (extract_model(response) || model).to_s
    resolved_provider = extract_provider(response)

    Metering::Recorder.record(
      model:         resolved_model,
      input_tokens:  usage[:input_tokens],
      output_tokens: usage[:output_tokens],
      provider:      resolved_provider
    )

    emit(
      provider:              resolved_provider,
      model_id:              resolved_model,
      request_type:          'chat',
      tier:                  extract_tier(response),
      input_tokens:          usage[:input_tokens],
      output_tokens:         usage[:output_tokens],
      thinking_tokens:       usage[:thinking_tokens] || 0,
      total_tokens:          usage[:input_tokens] + usage[:output_tokens],
      finish_reason:         extract_finish_reason(response),
      error_category:        extract_error_category(response),
      error_code:            extract_error_code(response),
      error_message:         extract_error_message(response),
      provider_instance:     extract_provider_instance(response),
      dispatch_path:         extract_dispatch_path(response),
      route_attempts:        extract_route_attempts(response),
      provider_response_ref: response.respond_to?(:provider_response_id) ? response.provider_response_id : nil,
      latency_ms:            extract_latency_ms(response),
      wall_clock_ms:         extract_wall_clock_ms(response),
      caller:                caller,
      event_type:            'llm_completion',
      status:                extract_status(response)
    )
    nil
  end
end

.load_transportObject



20
21
22
23
24
25
# File 'lib/legion/llm/metering.rb', line 20

def self.load_transport
  return unless defined?(Legion::Transport::Message)

  require_relative 'transport/exchanges/metering'
  require_relative 'transport/messages/metering_event'
end

.metering_event_classObject



146
147
148
149
150
151
152
153
154
155
156
# File 'lib/legion/llm/metering.rb', line 146

def metering_event_class
  return Legion::LLM::Transport::Messages::MeteringEvent if defined?(Legion::LLM::Transport::Messages::MeteringEvent)

  load_transport
  return Legion::LLM::Transport::Messages::MeteringEvent if defined?(Legion::LLM::Transport::Messages::MeteringEvent)

  Legion::LLM::Metering::Event
rescue NameError, LoadError => e
  handle_exception(e, level: :warn, handled: true, operation: 'llm.metering.event_class')
  nil
end

.read_spoolObject



236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
# File 'lib/legion/llm/metering.rb', line 236

def read_spool
  SPOOL_MUTEX.synchronize do
    path = spool_file_path
    return [] unless File.exist?(path)

    lines = File.readlines(path, chomp: true)
    lines.filter_map do |line|
      next if line.strip.empty?

      Legion::JSON.load(line)
    end
  end
rescue StandardError => e
  handle_exception(e, level: :warn, operation: 'llm.metering.read_spool')
  []
end

.spool_dir_pathObject



293
294
295
# File 'lib/legion/llm/metering.rb', line 293

def spool_dir_path
  File.dirname(spool_file_path)
end

.spool_event(event) ⇒ Object

--- Spool internals (private) ---



191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
# File 'lib/legion/llm/metering.rb', line 191

def spool_event(event)
  SPOOL_MUTEX.synchronize do
    ensure_spool_dir
    enforce_max_events
    json = Legion::JSON.dump(event)
    line = if encrypt_spool?
             Legion::Crypt.encrypt(json)
           else
             json
           end
    File.open(spool_file_path, 'a') { |f| f.puts(line) }
  end
  log.debug("[llm][metering] spool_event written provider=#{event[:provider]} model=#{event[:model_id]}")
rescue StandardError => e
  handle_exception(e, level: :warn, operation: 'llm.metering.spool_event')
end

.spool_file_pathObject

Resolve spool file path at call time, honouring operator-configured paths (e.g. for containerised deployments where $HOME is not writable). Falls back to the compile-time SPOOL_FILE constant.



288
289
290
291
# File 'lib/legion/llm/metering.rb', line 288

def spool_file_path
  configured = spool_settings[:path]
  configured && !configured.to_s.strip.empty? ? configured.to_s : SPOOL_FILE
end

.spool_settingsObject



280
281
282
283
# File 'lib/legion/llm/metering.rb', line 280

def spool_settings
  settings = Legion::Settings[:llm][:metering][:spool]
  settings.is_a?(Hash) ? settings : {}
end

.transport_connected?Boolean

Returns:

  • (Boolean)


142
143
144
# File 'lib/legion/llm/metering.rb', line 142

def transport_connected?
  Legion::Settings.dig(:transport, :connected) == true
end

.truncate_spoolObject



253
254
255
256
257
258
259
260
# File 'lib/legion/llm/metering.rb', line 253

def truncate_spool
  SPOOL_MUTEX.synchronize do
    path = spool_file_path
    File.write(path, '') if File.exist?(path)
  end
rescue StandardError => e
  handle_exception(e, level: :warn, operation: 'llm.metering.truncate_spool')
end