Class: Langfuse::Client

Inherits:
Object
  • Object
show all
Extended by:
SpanWrappers
Defined in:
lib/langfuse/client.rb

Defined Under Namespace

Classes: StdoutLogDevice

Constant Summary collapse

MAX_BATCH_SIZE_BYTES =

The ingestion API limits batch payloads to 3.5 MB in total

3_500_000
ENVIRONMENT_PATTERN =

Allowed format for the tracing environment field

/\A(?!langfuse)[a-z0-9\-_]{1,40}\z/
FLUSH_THREAD_JOIN_TIMEOUT =

How long shutdown waits for the flush thread to finish its current send

5
DROPPED_EVENTS_WARN_INTERVAL =

Dropped-event warnings are emitted on the first drop and then every N drops

100
INGESTION_MODES =

Supported ingestion transports: the legacy ingestion API and OTLP (Langfuse v4)

%i[legacy otel].freeze
RETRY_BASE_DELAY_SECONDS =

Backoff for transient request failures. The delay is jittered so that many clients hitting the same rate limit do not retry in lockstep, and capped so a server-sent Retry-After cannot stall a flush for minutes.

0.5
MAX_RETRY_DELAY_SECONDS =
10

Constants included from SpanWrappers

SpanWrappers::TYPES

Instance Attribute Summary collapse

Instance Method Summary collapse

Methods included from SpanWrappers

define_span_wrappers

Constructor Details

#initialize(public_key: nil, secret_key: nil, host: nil, debug: false, timeout: nil, retries: nil, flush_interval: nil, auto_flush: nil, ingestion_mode: nil, environment: nil, sample_rate: nil, mask: nil, flush_at: nil, max_queue_size: nil, logger: nil, shutdown_on_exit: nil, http_adapter: nil) ⇒ Client

timeout/retries default to nil so that Langfuse.configure values are not shadowed by the method defaults; the fallbacks live in config_value.



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
# File 'lib/langfuse/client.rb', line 55

def initialize(public_key: nil, secret_key: nil, host: nil, debug: false, timeout: nil, retries: nil,
               flush_interval: nil, auto_flush: nil, ingestion_mode: nil, environment: nil,
               sample_rate: nil, mask: nil, flush_at: nil, max_queue_size: nil, logger: nil,
               shutdown_on_exit: nil, http_adapter: nil)
  @public_key = config_value(public_key, 'LANGFUSE_PUBLIC_KEY', :public_key)
  @secret_key = config_value(secret_key, 'LANGFUSE_SECRET_KEY', :secret_key)
  @host = host || ENV['LANGFUSE_HOST'] || ENV['LANGFUSE_BASE_URL'] || Langfuse.configuration.host
  @debug = debug || ENV['LANGFUSE_DEBUG'] == 'true' || Langfuse.configuration.debug
  @timeout = config_value(timeout, nil, :timeout) { 30 }
  @retries = config_value(retries, nil, :retries) { 3 }
  @flush_interval = config_value(flush_interval, 'LANGFUSE_FLUSH_INTERVAL', :flush_interval) { 5 }
  @flush_at = config_value(flush_at, 'LANGFUSE_FLUSH_AT', :flush_at) { 15 }
  @max_queue_size = config_value(max_queue_size, 'LANGFUSE_MAX_QUEUE_SIZE', :max_queue_size) { 10_000 }
  @auto_flush = resolve_auto_flush(auto_flush)
  @logger = logger || Langfuse.configuration.logger || build_default_logger
  @ingestion_mode = resolve_ingestion_mode(ingestion_mode)
  @environment = resolve_environment(environment)
  @sample_rate = resolve_sample_rate(sample_rate)
  @mask = resolve_mask(mask)
  @shutdown_on_exit = shutdown_on_exit.nil? ? Langfuse.configuration.shutdown_on_exit : shutdown_on_exit
  @http_adapter = http_adapter || Langfuse.configuration.http_adapter
  @shutdown = false

  raise AuthenticationError, 'Public key is required' unless @public_key
  raise AuthenticationError, 'Secret key is required' unless @secret_key

  setup_transport
  @event_queue = Concurrent::Array.new
  @queue_mutex = Mutex.new
  @flush_mutex = Mutex.new
  @flush_condition = ConditionVariable.new
  @dropped_events = 0
  @prompt_cache = PromptCache.new
  start_flush_thread if @auto_flush
  register_shutdown_hook if @shutdown_on_exit
end

Instance Attribute Details

#auto_flush ⇒ Object (readonly)

Returns the value of attribute auto_flush.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def auto_flush
  @auto_flush
end

#debug ⇒ Object (readonly)

Returns the value of attribute debug.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def debug
  @debug
end

#environment ⇒ Object (readonly)

Returns the value of attribute environment.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def environment
  @environment
end

#flush_at ⇒ Object (readonly)

Returns the value of attribute flush_at.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def flush_at
  @flush_at
end

#flush_interval ⇒ Object (readonly)

Returns the value of attribute flush_interval.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def flush_interval
  @flush_interval
end

#host ⇒ Object (readonly)

Returns the value of attribute host.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def host
  @host
end

#ingestion_mode ⇒ Object (readonly)

Returns the value of attribute ingestion_mode.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def ingestion_mode
  @ingestion_mode
end

#logger ⇒ Object (readonly)

Returns the value of attribute logger.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def logger
  @logger
end

#mask ⇒ Object (readonly)

Returns the value of attribute mask.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def mask
  @mask
end

#max_queue_size ⇒ Object (readonly)

Returns the value of attribute max_queue_size.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def max_queue_size
  @max_queue_size
end

#public_key ⇒ Object (readonly)

Returns the value of attribute public_key.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def public_key
  @public_key
end

#retries ⇒ Object (readonly)

Returns the value of attribute retries.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def retries
  @retries
end

#sample_rate ⇒ Object (readonly)

Returns the value of attribute sample_rate.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def sample_rate
  @sample_rate
end

#secret_key ⇒ Object (readonly)

Returns the value of attribute secret_key.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def secret_key
  @secret_key
end

#timeout ⇒ Object (readonly)

Returns the value of attribute timeout.



50
51
52
# File 'lib/langfuse/client.rb', line 50

def timeout
  @timeout
end

Instance Method Details

#create_prompt(name:, prompt:, labels: [], config: {}, **kwargs) ⇒ Object



259
260
261
262
263
264
265
266
267
268
269
270
# File 'lib/langfuse/client.rb', line 259

def create_prompt(name:, prompt:, labels: [], config: {}, **kwargs)
  data = {
    name: name,
    prompt: prompt,
    labels: labels,
    config: config,
    **kwargs
  }

  response = post('/api/public/v2/prompts', data)
  Prompt.new(response.body)
end

#embedding(trace_id:, name: nil, start_time: nil, end_time: nil, input: nil, output: nil, model: nil, usage: nil, metadata: nil, level: nil, status_message: nil, parent_observation_id: nil, version: nil, **kwargs) ⇒ Object

Create an embedding observation (wrapper around span with as_type: 'embedding')



164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
# File 'lib/langfuse/client.rb', line 164

def embedding(trace_id:, name: nil, start_time: nil, end_time: nil, input: nil, output: nil,
              model: nil, usage: nil, metadata: nil, level: nil, status_message: nil,
              parent_observation_id: nil, version: nil, **kwargs)
   = ( || {}).merge(
    { model: model, usage: usage }.compact
  )
  span(
    trace_id: trace_id,
    name: name,
    start_time: start_time,
    end_time: end_time,
    input: input,
    output: output,
    metadata: .empty? ? nil : ,
    level: level,
    status_message: status_message,
    parent_observation_id: parent_observation_id,
    version: version,
    as_type: ObservationType::EMBEDDING,
    **kwargs
  )
end

#enqueue_event(type, body, trace_ref: nil) ⇒ Object

Event queue management



304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
# File 'lib/langfuse/client.rb', line 304

def enqueue_event(type, body, trace_ref: nil)
  # 验证事件类型是否有效
  valid_types = %w[
    trace-create trace-update
    generation-create generation-update
    span-create span-update
    event-create
    score-create
  ]

  unless valid_types.include?(type)
    @logger.debug { "Warning: Invalid event type '#{type}'. Skipping event." }
    return
  end

  # Runs before the event is queued so events inherited from a parent process
  # are discarded without dropping the event we are about to enqueue.
  ensure_flush_thread

  prepared_body = prepare_queued_body(body)

  return unless sampled_event?(type, prepared_body)

  event = {
    id: Utils.generate_id,
    type: type,
    timestamp: Utils.current_timestamp,
    body: prepared_body
  }
  event[:trace_ref] = trace_ref if trace_ref

  # The queue is drained under the same lock, so a concurrent flush can no
  # longer take an event out between finding it and merging into it.
  queued = @queue_mutex.synchronize do
    if type == 'trace-update'
      merge_or_queue_trace_update?(event)
    else
      push_event?(event)
    end
  end
  return unless queued

  @logger.debug { "Enqueued event: #{type}" }

  request_flush if @auto_flush && @event_queue.length >= @flush_at
end

#event(trace_id:, name:, id: nil, start_time: nil, input: nil, output: nil, metadata: nil, level: nil, status_message: nil, parent_observation_id: nil, version: nil, **kwargs) ⇒ Object

Event operations



219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
# File 'lib/langfuse/client.rb', line 219

def event(trace_id:, name:, id: nil, start_time: nil, input: nil, output: nil, metadata: nil,
          level: nil, status_message: nil, parent_observation_id: nil, version: nil, **kwargs)
  Event.new(
    client: self,
    trace_id: trace_id,
    id: id || generate_observation_id,
    name: name,
    start_time: start_time,
    input: input,
    output: output,
    metadata: ,
    level: level,
    status_message: status_message,
    parent_observation_id: parent_observation_id,
    version: version,
    **kwargs
  )
end

#flush ⇒ Object



351
352
353
354
355
356
357
358
# File 'lib/langfuse/client.rb', line 351

def flush
  events = @queue_mutex.synchronize do
    @event_queue.empty? ? [] : @event_queue.shift(@event_queue.length)
  end
  return if events.empty?

  send_batch(events)
end

#generate_observation_id ⇒ Object

Generate an observation ID matching the active ingestion mode (W3C 16-char hex for :otel, UUID for :legacy)



106
107
108
# File 'lib/langfuse/client.rb', line 106

def generate_observation_id
  @ingestion_mode == :otel ? Utils.generate_hex_span_id : Utils.generate_id
end

#generate_trace_id ⇒ Object

Generate a trace ID matching the active ingestion mode (W3C 32-char hex for :otel, UUID for :legacy)



100
101
102
# File 'lib/langfuse/client.rb', line 100

def generate_trace_id
  @ingestion_mode == :otel ? Utils.generate_hex_trace_id : Utils.generate_id
end

#generation(trace_id:, id: nil, name: nil, start_time: nil, end_time: nil, completion_start_time: nil, model: nil, model_parameters: nil, input: nil, output: nil, usage: nil, usage_details: nil, cost_details: nil, prompt: nil, metadata: nil, level: nil, status_message: nil, parent_observation_id: nil, version: nil, **kwargs) ⇒ Object

Generation operations



188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
# File 'lib/langfuse/client.rb', line 188

def generation(trace_id:, id: nil, name: nil, start_time: nil, end_time: nil, completion_start_time: nil,
               model: nil, model_parameters: nil, input: nil, output: nil, usage: nil,
               usage_details: nil, cost_details: nil, prompt: nil,
               metadata: nil, level: nil, status_message: nil, parent_observation_id: nil,
               version: nil, **kwargs)
  Generation.new(
    client: self,
    trace_id: trace_id,
    id: id || generate_observation_id,
    name: name,
    start_time: start_time || Utils.current_timestamp,
    end_time: end_time,
    completion_start_time: completion_start_time,
    model: model,
    model_parameters: model_parameters,
    input: input,
    output: output,
    usage: usage,
    usage_details: usage_details,
    cost_details: cost_details,
    prompt: prompt,
    metadata: ,
    level: level,
    status_message: status_message,
    parent_observation_id: parent_observation_id,
    version: version,
    **kwargs
  )
end

#get_prompt(name, version: nil, label: nil, cache_ttl_seconds: 60, retries: nil) ⇒ Object

Prompt operations



239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
# File 'lib/langfuse/client.rb', line 239

def get_prompt(name, version: nil, label: nil, cache_ttl_seconds: 60, retries: nil)
  cache_key = "prompt:#{name}:#{version}:#{label}"
  cached = @prompt_cache.read(cache_key, cache_ttl_seconds)
  return cached if cached

  begin
    prompt = request_prompt(name, version: version, label: label, retries: retries)
  rescue StandardError => e
    # An expired entry beats no prompt at all: serving it keeps the
    # application running through a Langfuse outage.
    stale = @prompt_cache.read_stale(cache_key)
    raise unless stale

    @logger.warn("Langfuse prompt fetch failed (#{name}), serving the cached copy: #{e.message}")
    return stale
  end

  @prompt_cache.write(cache_key, prompt)
end

#inspect ⇒ Object

Keep the secret key out of logs, console output and exception messages.



93
94
95
96
# File 'lib/langfuse/client.rb', line 93

def inspect
  "#<#{self.class.name} host=#{@host.inspect} public_key=#{@public_key.inspect} " \
    "ingestion_mode=#{@ingestion_mode.inspect}>"
end

#score(name:, value:, trace_id: nil, observation_id: nil, session_id: nil, dataset_run_id: nil, id: nil, data_type: nil, comment: nil, metadata: nil, config_id: nil, queue_id: nil, environment: nil, **kwargs) ⇒ Object Also known as: create_score

Score/Evaluation operations Scores can target a trace, an observation (trace_id + observation_id), a session (session_id) or a dataset run (dataset_run_id).



275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
# File 'lib/langfuse/client.rb', line 275

def score(name:, value:, trace_id: nil, observation_id: nil, session_id: nil, dataset_run_id: nil,
          id: nil, data_type: nil, comment: nil, metadata: nil, config_id: nil, queue_id: nil,
          environment: nil, **kwargs)
  data = {
    id: id,
    trace_id: trace_id,
    observation_id: observation_id,
    session_id: session_id,
    dataset_run_id: dataset_run_id,
    name: name,
    value: value,
    data_type: data_type,
    comment: comment,
    metadata: ,
    config_id: config_id,
    queue_id: queue_id,
    environment: environment,
    **kwargs
  }.compact

  if trace_id.nil? && observation_id.nil? && session_id.nil? && dataset_run_id.nil?
    @logger.warn('Langfuse score should reference a trace_id, observation_id, session_id or dataset_run_id')
  end

  enqueue_event('score-create', data)
end

#shutdown ⇒ Object



360
361
362
363
364
365
366
# File 'lib/langfuse/client.rb', line 360

def shutdown
  return if @shutdown

  @shutdown = true
  stop_flush_thread
  flush unless @event_queue.empty?
end

#span(trace_id:, id: nil, name: nil, start_time: nil, end_time: nil, input: nil, output: nil, metadata: nil, level: nil, status_message: nil, parent_observation_id: nil, version: nil, as_type: nil, **kwargs) ⇒ Object

Span operations



131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
# File 'lib/langfuse/client.rb', line 131

def span(trace_id:, id: nil, name: nil, start_time: nil, end_time: nil, input: nil, output: nil,
         metadata: nil, level: nil, status_message: nil, parent_observation_id: nil,
         version: nil, as_type: nil, **kwargs)
  Span.new(
    client: self,
    trace_id: trace_id,
    id: id || generate_observation_id,
    name: name,
    start_time: start_time || Utils.current_timestamp,
    end_time: end_time,
    input: input,
    output: output,
    metadata: ,
    level: level,
    status_message: status_message,
    parent_observation_id: parent_observation_id,
    version: version,
    as_type: as_type,
    **kwargs
  )
end

#trace(id: nil, name: nil, user_id: nil, session_id: nil, version: nil, release: nil, input: nil, output: nil, metadata: nil, tags: nil, timestamp: nil, **kwargs) ⇒ Object

Trace operations



111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
# File 'lib/langfuse/client.rb', line 111

def trace(id: nil, name: nil, user_id: nil, session_id: nil, version: nil, release: nil,
          input: nil, output: nil, metadata: nil, tags: nil, timestamp: nil, **kwargs)
  Trace.new(
    client: self,
    id: id || generate_trace_id,
    name: name,
    user_id: user_id,
    session_id: session_id,
    version: version,
    release: release,
    input: input,
    output: output,
    metadata: ,
    tags: tags,
    timestamp: timestamp || Utils.current_timestamp,
    **kwargs
  )
end