Class: Langfuse::Client
- Inherits:
-
Object
- Object
- Langfuse::Client
- 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
Instance Attribute Summary collapse
-
#auto_flush ⇒ Object
readonly
Returns the value of attribute auto_flush.
-
#debug ⇒ Object
readonly
Returns the value of attribute debug.
-
#environment ⇒ Object
readonly
Returns the value of attribute environment.
-
#flush_at ⇒ Object
readonly
Returns the value of attribute flush_at.
-
#flush_interval ⇒ Object
readonly
Returns the value of attribute flush_interval.
-
#host ⇒ Object
readonly
Returns the value of attribute host.
-
#ingestion_mode ⇒ Object
readonly
Returns the value of attribute ingestion_mode.
-
#logger ⇒ Object
readonly
Returns the value of attribute logger.
-
#mask ⇒ Object
readonly
Returns the value of attribute mask.
-
#max_queue_size ⇒ Object
readonly
Returns the value of attribute max_queue_size.
-
#public_key ⇒ Object
readonly
Returns the value of attribute public_key.
-
#retries ⇒ Object
readonly
Returns the value of attribute retries.
-
#sample_rate ⇒ Object
readonly
Returns the value of attribute sample_rate.
-
#secret_key ⇒ Object
readonly
Returns the value of attribute secret_key.
-
#timeout ⇒ Object
readonly
Returns the value of attribute timeout.
Instance Method Summary collapse
- #create_prompt(name:, prompt:, labels: [], config: {}, **kwargs) ⇒ Object
-
#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').
-
#enqueue_event(type, body, trace_ref: nil) ⇒ Object
Event queue management.
-
#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.
- #flush ⇒ Object
-
#generate_observation_id ⇒ Object
Generate an observation ID matching the active ingestion mode (W3C 16-char hex for :otel, UUID for :legacy).
-
#generate_trace_id ⇒ Object
Generate a trace ID matching the active ingestion mode (W3C 32-char hex for :otel, UUID for :legacy).
-
#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.
-
#get_prompt(name, version: nil, label: nil, cache_ttl_seconds: 60, retries: nil) ⇒ Object
Prompt operations.
-
#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
constructor
timeout/retries default to nil so that Langfuse.configure values are not shadowed by the method defaults; the fallbacks live in config_value.
-
#inspect ⇒ Object
Keep the secret key out of logs, console output and exception messages.
-
#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: #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).
- #shutdown ⇒ Object
-
#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.
-
#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.
Methods included from SpanWrappers
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 (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: , 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., 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: , 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., 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: , 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.}") 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., end_time: end_time, input: input, output: output, metadata: , level: level, 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: , timestamp: || Utils., **kwargs ) end |