Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Fixed
- Emit Langfuse v4 experiment attributes for local and dataset-backed experiment observations (#118).

## [0.12.0] - 2026-08-20

### Changed
Expand Down
2 changes: 2 additions & 0 deletions docs/API_REFERENCE.md
Original file line number Diff line number Diff line change
Expand Up @@ -1547,6 +1547,8 @@ run_experiment(name:, task:, data: nil, dataset_name: nil, description: nil,

**Returns:** `ExperimentResult`

The result exposes `experiment_id`, which is the dataset run ID for a dataset-backed experiment or a generated run-scoped ID for local data. See [EXPERIMENTS.md](EXPERIMENTS.md#langfuse-v4-attribution) for the observation attributes and propagation behavior.

**Raises:** `ArgumentError` if both or neither of `data`/`dataset_name` provided

**Example:**
Expand Down
3 changes: 3 additions & 0 deletions docs/DATASETS.md
Original file line number Diff line number Diff line change
Expand Up @@ -193,6 +193,9 @@ end

The block receives a traced span. On completion (or error), the trace is flushed and the item is linked automatically.

> [!NOTE]
> Manual `item.link` calls and the lower-level `item.run` helper create dataset run links, but they cannot add Langfuse v4 experiment attributes to observations that have already started. Use `run_experiment` when the run must appear as a v4 experiment with complete observation attribution.

## Managing Dataset Runs

Dataset runs are created when you link items into a named run, either manually or through experiments. Once a run exists, you can fetch it, list all runs for a dataset, or delete it by dataset name and run name.
Expand Down
10 changes: 10 additions & 0 deletions docs/EXPERIMENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ result = client.run_experiment(

puts result.format
puts "#{result.successes.size} passed, #{result.failures.size} failed"
puts result.experiment_id # => experiment identifier shared by the run
puts result.dataset_run_url # => link to Langfuse UI
```

Expand Down Expand Up @@ -61,6 +62,14 @@ result = client.run_experiment(

Each hash is wrapped into an `ExperimentItem` struct with `input`, `expected_output`, and `metadata` fields. Both symbol and string keys are accepted.

## Langfuse v4 Attribution

`run_experiment` adds the attributes that Langfuse v4 uses to discover experiments. The item root observation and every child observation share the experiment ID, run name, item ID, item root observation ID, experiment metadata, item metadata, and `sdk-experiment` environment.

The run description and expected output stay on the item root observation. Item-level scores target both the trace and the item root observation.

For a dataset-backed experiment, `result.experiment_id` is the server-provided dataset run ID when the first link succeeds. If the run selects a generated fallback ID after a link failure, all later items keep that ID, and `dataset_run_id` and `dataset_run_url` remain `nil`. A local-data experiment receives one generated 16-character hexadecimal ID shared by all items in that run. Local item IDs are the first 16 hexadecimal characters of a SHA-256 digest over the SDK-serialized input.

## Parameters

### `Client#run_experiment`
Expand Down Expand Up @@ -174,6 +183,7 @@ Returned by `run_experiment`.
| `description` | String, nil | Run description |
| `item_results` | Array\<ItemResult\> | All per-item results |
| `run_evaluations` | Array\<Evaluation\> | Run-level evaluation results |
| `experiment_id` | String, nil | Experiment ID shared by the run |
| `dataset_run_id` | String, nil | Dataset run ID from the server |
| `dataset_run_url` | String, nil | URL to the run in Langfuse UI |

Expand Down
1 change: 1 addition & 0 deletions lib/langfuse.rb
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,7 @@ class UnauthorizedError < ApiError; end
require_relative "langfuse/masking"
require_relative "langfuse/masking_exporter"
require_relative "langfuse/otel_attributes"
require_relative "langfuse/experiment_attributes"
require_relative "langfuse/propagation"
require_relative "langfuse/app_root_tracking"
require_relative "langfuse/span_processor"
Expand Down
1 change: 1 addition & 0 deletions lib/langfuse/dataset_item_client.rb
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,7 @@ def link(trace_id:, run_name:, observation_id: nil, metadata: nil, run_descripti
#
# Executes the block inside an observed span, flushes the trace, then
# creates a dataset run item linking this item to the resulting trace.
# Use {Client#run_experiment} when Langfuse v4 experiment attribution is required.
#
# @param run_name [String] run name for grouping
# @param run_description [String, nil] optional run description
Expand Down
67 changes: 67 additions & 0 deletions lib/langfuse/experiment_attributes.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
# frozen_string_literal: true

require "digest"
require "securerandom"

module Langfuse
# Builds the OpenTelemetry attributes that identify v4 experiment observations.
#
# @api private
module ExperimentAttributes
ENVIRONMENT = "sdk-experiment"

def self.generate_experiment_id
SecureRandom.hex(8)
end

def self.generate_item_id(input)
serialized_input = OtelAttributes.serialize(input, preserve_strings: true) || input.to_s
Digest::SHA256.hexdigest(serialized_input)[0, 16]
end

# rubocop:disable Metrics/ParameterLists
def self.propagated(experiment_id:, run_name:, item_id:, root_observation_id:,
dataset_id: nil, experiment_metadata: nil, item_metadata: nil, mask: nil)
attributes = {
OtelAttributes::EXPERIMENT_ID => experiment_id,
OtelAttributes::EXPERIMENT_NAME => run_name,
OtelAttributes::EXPERIMENT_DATASET_ID => dataset_id,
OtelAttributes::EXPERIMENT_ITEM_ID => item_id,
OtelAttributes::EXPERIMENT_ITEM_ROOT_OBSERVATION_ID => root_observation_id,
OtelAttributes::ENVIRONMENT => ENVIRONMENT
}.compact
attributes.merge!(masked_metadata(experiment_metadata, OtelAttributes::EXPERIMENT_METADATA, mask))
attributes.merge!(masked_metadata(item_metadata, OtelAttributes::EXPERIMENT_ITEM_METADATA, mask))
end
# rubocop:enable Metrics/ParameterLists

def self.root(description:, expected_output:, mask: nil)
masked_output = Masking.apply(expected_output, mask: mask)
{
OtelAttributes::EXPERIMENT_DESCRIPTION => description,
OtelAttributes::EXPERIMENT_ITEM_EXPECTED_OUTPUT =>
OtelAttributes.serialize(masked_output, preserve_strings: true)
}.compact
end

def self.observation_metadata(name:, run_name:, experiment_metadata:, item_metadata:,
dataset_id: nil, dataset_item_id: nil)
metadata = hash_metadata(item_metadata).merge(hash_metadata(experiment_metadata))
metadata[:experiment_name] = name
metadata[:experiment_run_name] = run_name
metadata[:dataset_id] = dataset_id if dataset_id
metadata[:dataset_item_id] = dataset_item_id if dataset_item_id
metadata
end

def self.masked_metadata(metadata, prefix, mask)
OtelAttributes.flatten_metadata(Masking.apply(metadata, mask: mask), prefix)
end
private_class_method :masked_metadata

def self.hash_metadata(metadata)
metadata.is_a?(Hash) ? metadata : {}
end
private_class_method :hash_metadata
end
end
7 changes: 5 additions & 2 deletions lib/langfuse/experiment_result.rb
Original file line number Diff line number Diff line change
Expand Up @@ -18,26 +18,29 @@ class ExperimentResult
# @return [String, nil] run description
# @return [Array<ItemResult>] per-item results (all items, including failures)
# @return [Array<Evaluation>] run-level evaluation results
# @return [String, nil] resolved experiment ID shared by the run
# @return [String, nil] dataset run ID from the server
# @return [String, nil] URL to the dataset run in Langfuse UI
attr_reader :name, :run_name, :description, :item_results, :run_evaluations,
:dataset_run_id, :dataset_run_url
:experiment_id, :dataset_run_id, :dataset_run_url

# @param name [String] experiment/run name
# @param item_results [Array<ItemResult>] per-item results
# @param run_evaluations [Array<Evaluation>] run-level evaluations
# @param run_name [String, nil] auto-generated run name
# @param description [String, nil] run description
# @param experiment_id [String, nil] resolved experiment ID shared by the run
# @param dataset_run_id [String, nil] dataset run ID from the server
# @param dataset_run_url [String, nil] URL to the dataset run in Langfuse UI
# rubocop:disable Metrics/ParameterLists
def initialize(name:, item_results:, run_evaluations: [], run_name: nil, description: nil,
dataset_run_id: nil, dataset_run_url: nil)
experiment_id: nil, dataset_run_id: nil, dataset_run_url: nil)
@name = name
@item_results = item_results
@run_evaluations = run_evaluations
@run_name = run_name
@description = description
@experiment_id = experiment_id
@dataset_run_id = dataset_run_id
@dataset_run_url = dataset_run_url
end
Expand Down
93 changes: 78 additions & 15 deletions lib/langfuse/experiment_runner.rb
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ def initialize(client:, name:, items:, task:, evaluators: [], run_evaluators: []
@description = description
@run_name = run_name || "#{name} - #{Time.now.utc.iso8601}"
@logger = Langfuse.configuration.logger
@fallback_experiment_id = ExperimentAttributes.generate_experiment_id
@resolved_experiment_id = nil
@dataset_run_id = nil
@dataset_id = nil
end
Expand All @@ -53,6 +55,7 @@ def execute
run_evaluations: run_evals,
run_name: @run_name,
description: @description,
experiment_id: resolved_experiment_id,
dataset_run_id: @dataset_run_id,
dataset_run_url: build_dataset_run_url
)
Expand All @@ -68,7 +71,7 @@ def process_item(item, index)
return ItemResult.new(item: item, trace_id: trace_id, observation_id: observation_id, error: task_error)
end

evaluations = execute_evaluators(item, output, trace_id)
evaluations = execute_evaluators(item, output, trace_id, observation_id)
ItemResult.new(item: item, output: output, trace_id: trace_id,
observation_id: observation_id, evaluations: evaluations)
end
Expand All @@ -77,21 +80,19 @@ def run_task_in_trace(item)
TracedExecution.call(
trace_name: "experiment-#{@name}",
input: item.input,
metadata: @metadata,
task: ->(_span) { @task.call(item) }
) do |span, trace_id|
# Link before running task — server accepts forward-referenced trace IDs
link_to_dataset_run(item, trace_id, span.id) if item.is_a?(DatasetItemClient)
end
metadata: observation_metadata(item),
task: ->(_span) { @task.call(item) },
prepare_context: ->(span, trace_id) { prepare_experiment_context(item, span, trace_id) }
)
end

def execute_evaluators(item, output, trace_id)
def execute_evaluators(item, output, trace_id, observation_id)
evaluations = @evaluators.flat_map do |evaluator|
raw_result = call_evaluator(evaluator, item, output)
normalize_evaluator_result(raw_result, source: "Evaluator")
end

evaluations.each { |evaluation| persist_score(evaluation, trace_id) }
evaluations.each { |evaluation| persist_score(evaluation, trace_id, observation_id) }
evaluations
end

Expand All @@ -104,10 +105,11 @@ def call_evaluator(evaluator, item, output)
nil
end

def persist_score(evaluation, trace_id)
def persist_score(evaluation, trace_id, observation_id)
@client.create_score(
name: evaluation.name, value: evaluation.value,
trace_id: trace_id, comment: evaluation.comment, data_type: evaluation.data_type,
trace_id: trace_id, observation_id: observation_id,
comment: evaluation.comment, data_type: evaluation.data_type,
config_id: evaluation.config_id, metadata: evaluation.metadata
)
rescue StandardError => e
Expand Down Expand Up @@ -145,13 +147,74 @@ def link_to_dataset_run(item, trace_id, observation_id)
trace_id: trace_id, observation_id: observation_id,
metadata: @metadata, run_description: @description
)
unless @dataset_run_id
@dataset_run_id = response&.dig("datasetRunId")
@dataset_id = item.dataset_id if @dataset_run_id
end
capture_dataset_run_identity(item, response)
response
rescue StandardError => e
@logger.warn("Dataset run item linking failed: #{e.message}")
nil
end

# Do not expose a server run ID that differs from the identity already attached to observations.
def capture_dataset_run_identity(item, response)
candidate_id = response&.dig("datasetRunId")
return unless candidate_id
return if @resolved_experiment_id && @resolved_experiment_id != candidate_id

@dataset_run_id = candidate_id if @dataset_run_id.nil?
@dataset_id = item.dataset_id if @dataset_id.nil?
end

def prepare_experiment_context(item, span, trace_id)
response = link_to_dataset_run(item, trace_id, span.id) if item.is_a?(DatasetItemClient)
root_experiment_attributes(item).each do |key, value|
span.otel_span.set_attribute(key, value)
end
propagated_experiment_attributes(item, span.id, response)
end

def root_experiment_attributes(item)
ExperimentAttributes.root(
description: @description,
expected_output: item.expected_output,
mask: Langfuse.configuration.mask
)
end

def propagated_experiment_attributes(item, observation_id, response)
ExperimentAttributes.propagated(
experiment_id: resolved_experiment_id(response),
run_name: @run_name,
dataset_id: item_dataset_id(item),
item_id: item_id(item),
root_observation_id: observation_id,
experiment_metadata: @metadata,
item_metadata: item_metadata(item),
mask: Langfuse.configuration.mask
)
end

# A run keeps its first resolved identity even if later dataset links have a different result.
def resolved_experiment_id(response = nil)
@resolved_experiment_id ||= response&.dig("datasetRunId") || @fallback_experiment_id
end

def observation_metadata(item)
ExperimentAttributes.observation_metadata(
name: @name,
run_name: @run_name,
experiment_metadata: @metadata,
item_metadata: item_metadata(item),
dataset_id: item_dataset_id(item),
dataset_item_id: item.is_a?(DatasetItemClient) ? item.id : nil
)
end

def item_id(item)
item.is_a?(DatasetItemClient) ? item.id : ExperimentAttributes.generate_item_id(item.input)
end

def item_dataset_id(item)
item.dataset_id if item.is_a?(DatasetItemClient)
end

def flush_all
Expand Down
11 changes: 11 additions & 0 deletions lib/langfuse/otel_attributes.rb
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,17 @@ module OtelAttributes
OBSERVATION_COMPLETION_START_TIME = "langfuse.observation.completion_start_time"
IS_APP_ROOT = "langfuse.internal.is_app_root"

# Experiment attributes
EXPERIMENT_ID = "langfuse.experiment.id"
EXPERIMENT_NAME = "langfuse.experiment.name"
EXPERIMENT_DESCRIPTION = "langfuse.experiment.description"
EXPERIMENT_METADATA = "langfuse.experiment.metadata"
EXPERIMENT_DATASET_ID = "langfuse.experiment.dataset.id"
EXPERIMENT_ITEM_ID = "langfuse.experiment.item.id"
EXPERIMENT_ITEM_EXPECTED_OUTPUT = "langfuse.experiment.item.expected_output"
EXPERIMENT_ITEM_METADATA = "langfuse.experiment.item.metadata"
EXPERIMENT_ITEM_ROOT_OBSERVATION_ID = "langfuse.experiment.item.root_observation_id"

# Common attributes
VERSION = "langfuse.version"
RELEASE = "langfuse.release"
Expand Down
25 changes: 24 additions & 1 deletion lib/langfuse/propagation.rb
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,11 @@ module Propagation
ENVIRONMENT_VALUE_PATTERN = /\A(?!langfuse)[a-z0-9_-]+\z/
private_constant :ENVIRONMENT_VALUE_PATTERN

EXPERIMENT_ATTRIBUTES_CONTEXT_KEY = OpenTelemetry::Context.create_key(
"#{BAGGAGE_PREFIX}experiment_attributes"
)
private_constant :EXPERIMENT_ATTRIBUTES_CONTEXT_KEY

# Map of propagated attribute keys to span attribute keys
SPAN_KEY_MAP = {
"user_id" => OtelAttributes::TRACE_USER_ID,
Expand Down Expand Up @@ -186,7 +191,25 @@ def self.get_propagated_attributes_from_context(context)
end
# rubocop:enable Metrics/AbcSize, Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity

propagated_attributes
experiment_attributes = context.value(EXPERIMENT_ATTRIBUTES_CONTEXT_KEY) || {}
propagated_attributes.merge(experiment_attributes)
end

# Apply SDK-owned experiment attributes to the current span and future children.
#
# @param attributes [Hash<String, Object>] serialized OpenTelemetry experiment attributes
# @yield Block within which experiment attributes propagate
# @return [Object] the block result
# @api private
def self._with_experiment_attributes(attributes, &)
return yield if attributes.nil? || attributes.empty?

frozen_attributes = attributes.dup.freeze
current_span = OpenTelemetry::Trace.current_span
frozen_attributes.each { |key, value| current_span.set_attribute(key, value) } if current_span.recording?

context = OpenTelemetry::Context.current.set_value(EXPERIMENT_ATTRIBUTES_CONTEXT_KEY, frozen_attributes)
OpenTelemetry::Context.with_current(context, &)
end

# Merge metadata with existing context value
Expand Down
Loading
Loading