diff --git a/CHANGELOG.md b/CHANGELOG.md index e3e2f34..ba93f79 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/docs/API_REFERENCE.md b/docs/API_REFERENCE.md index e9db91c..a615156 100644 --- a/docs/API_REFERENCE.md +++ b/docs/API_REFERENCE.md @@ -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:** diff --git a/docs/DATASETS.md b/docs/DATASETS.md index 745d109..de5e3c5 100644 --- a/docs/DATASETS.md +++ b/docs/DATASETS.md @@ -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. diff --git a/docs/EXPERIMENTS.md b/docs/EXPERIMENTS.md index 10d9c81..2062304 100644 --- a/docs/EXPERIMENTS.md +++ b/docs/EXPERIMENTS.md @@ -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 ``` @@ -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` @@ -174,6 +183,7 @@ Returned by `run_experiment`. | `description` | String, nil | Run description | | `item_results` | Array\ | All per-item results | | `run_evaluations` | Array\ | 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 | diff --git a/lib/langfuse.rb b/lib/langfuse.rb index b6cb6b3..3e4aa0d 100644 --- a/lib/langfuse.rb +++ b/lib/langfuse.rb @@ -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" diff --git a/lib/langfuse/dataset_item_client.rb b/lib/langfuse/dataset_item_client.rb index c8df97a..f1aca45 100644 --- a/lib/langfuse/dataset_item_client.rb +++ b/lib/langfuse/dataset_item_client.rb @@ -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 diff --git a/lib/langfuse/experiment_attributes.rb b/lib/langfuse/experiment_attributes.rb new file mode 100644 index 0000000..7f133c1 --- /dev/null +++ b/lib/langfuse/experiment_attributes.rb @@ -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 diff --git a/lib/langfuse/experiment_result.rb b/lib/langfuse/experiment_result.rb index b632696..8413157 100644 --- a/lib/langfuse/experiment_result.rb +++ b/lib/langfuse/experiment_result.rb @@ -18,26 +18,29 @@ class ExperimentResult # @return [String, nil] run description # @return [Array] per-item results (all items, including failures) # @return [Array] 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] per-item results # @param run_evaluations [Array] 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 diff --git a/lib/langfuse/experiment_runner.rb b/lib/langfuse/experiment_runner.rb index c931d97..dcdf529 100644 --- a/lib/langfuse/experiment_runner.rb +++ b/lib/langfuse/experiment_runner.rb @@ -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 @@ -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 ) @@ -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 @@ -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 @@ -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 @@ -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 diff --git a/lib/langfuse/otel_attributes.rb b/lib/langfuse/otel_attributes.rb index 110bdd7..84faeb2 100644 --- a/lib/langfuse/otel_attributes.rb +++ b/lib/langfuse/otel_attributes.rb @@ -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" diff --git a/lib/langfuse/propagation.rb b/lib/langfuse/propagation.rb index a2ded91..953706d 100644 --- a/lib/langfuse/propagation.rb +++ b/lib/langfuse/propagation.rb @@ -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, @@ -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] 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 diff --git a/lib/langfuse/traced_execution.rb b/lib/langfuse/traced_execution.rb index 6b6dccc..1088b39 100644 --- a/lib/langfuse/traced_execution.rb +++ b/lib/langfuse/traced_execution.rb @@ -15,9 +15,9 @@ module TracedExecution # @param input [Object] input set on the root observation # @param metadata [Hash] metadata set on the root observation and trace # @param task [Proc] the callable to execute — receives the span - # @yield [span, trace_id] optional pre-task hook (e.g., dataset run linking) + # @param prepare_context [#call, nil] optional callable returning SDK-owned propagated attributes # @return [Array<(Object, String, String, StandardError | nil)>] output, trace_id, observation_id, error - def self.call(trace_name:, input:, task:, metadata: {}) + def self.call(trace_name:, input:, task:, metadata: {}, prepare_context: nil) trace_id = nil observation_id = nil output = nil @@ -26,9 +26,11 @@ def self.call(trace_name:, input:, task:, metadata: {}) Langfuse.observe(trace_name, input: input, metadata: metadata) do |span| trace_id = span.trace_id observation_id = span.id - Langfuse.propagate_attributes(trace_name: trace_name, metadata: metadata) do - yield(span, trace_id) if block_given? - output, task_error = execute_task(span, task) + experiment_attributes = prepare_context&.call(span, trace_id) || {} + Propagation._with_experiment_attributes(experiment_attributes) do + Langfuse.propagate_attributes(trace_name: trace_name, metadata: metadata) do + output, task_error = execute_task(span, task) + end end end diff --git a/spec/langfuse/experiment_attributes_spec.rb b/spec/langfuse/experiment_attributes_spec.rb new file mode 100644 index 0000000..9c6b3e4 --- /dev/null +++ b/spec/langfuse/experiment_attributes_spec.rb @@ -0,0 +1,114 @@ +# frozen_string_literal: true + +require "digest" + +RSpec.describe Langfuse::ExperimentAttributes do + describe ".generate_experiment_id" do + it "returns unique 16-character hexadecimal identifiers" do + first_id = described_class.generate_experiment_id + second_id = described_class.generate_experiment_id + + expect(first_id).to match(/\A[0-9a-f]{16}\z/) + expect(second_id).to match(/\A[0-9a-f]{16}\z/) + expect(second_id).not_to eq(first_id) + end + end + + describe ".generate_item_id" do + it "hashes the SDK-serialized input deterministically" do + input = { question: "What is 2 + 2?" } + serialized = Langfuse::OtelAttributes.serialize(input, preserve_strings: true) + expected_id = Digest::SHA256.hexdigest(serialized)[0, 16] + + expect(described_class.generate_item_id(input)).to eq(expected_id) + expect(described_class.generate_item_id(input)).to eq(expected_id) + end + + it "falls back to the input string when serialization fails" do + input = Object.new + allow(input).to receive(:to_json).and_raise(StandardError, "cannot serialize") + allow(input).to receive(:to_s).and_return("stable input") + + expected_id = Digest::SHA256.hexdigest("stable input")[0, 16] + expect(described_class.generate_item_id(input)).to eq(expected_id) + end + end + + describe ".propagated" do + it "builds v4 identity, dataset, metadata, and environment attributes" do + attributes = described_class.propagated( + experiment_id: "run-id", run_name: "nightly", + dataset_id: "dataset-id", item_id: "item-id", root_observation_id: "observation-id", + experiment_metadata: { model: { name: "test" } }, item_metadata: { difficulty: "easy" } + ) + + expect(attributes).to include( + "langfuse.experiment.id" => "run-id", + "langfuse.experiment.name" => "nightly", + "langfuse.experiment.dataset.id" => "dataset-id", + "langfuse.experiment.item.id" => "item-id", + "langfuse.experiment.item.root_observation_id" => "observation-id", + "langfuse.experiment.metadata.model.name" => "test", + "langfuse.experiment.item.metadata.difficulty" => "easy", + "langfuse.environment" => "sdk-experiment" + ) + end + + it "applies masking before flattening experiment and item metadata" do + mask = lambda do |data:| + data.transform_values { |value| value == "secret" ? "redacted" : value } + end + + attributes = described_class.propagated( + experiment_id: "run-id", run_name: "nightly", item_id: "item-id", + root_observation_id: "observation-id", experiment_metadata: { token: "secret" }, + item_metadata: { answer: "secret" }, mask: mask + ) + + expect(attributes["langfuse.experiment.metadata.token"]).to eq("redacted") + expect(attributes["langfuse.experiment.item.metadata.answer"]).to eq("redacted") + end + end + + describe ".root" do + it "keeps description and false expected output on the item root" do + attributes = described_class.root(description: "test run", expected_output: false) + + expect(attributes).to eq( + "langfuse.experiment.description" => "test run", + "langfuse.experiment.item.expected_output" => "false" + ) + end + + it "omits nil expected output" do + attributes = described_class.root(description: nil, expected_output: nil) + + expect(attributes).to eq({}) + end + + it "masks expected output before serialization" do + mask = ->(data:) { data == "secret" ? "redacted" : data } + + attributes = described_class.root(description: nil, expected_output: "secret", mask: mask) + + expect(attributes["langfuse.experiment.item.expected_output"]).to eq("redacted") + end + end + + describe ".observation_metadata" do + it "merges item metadata before run metadata and adds fixed experiment fields" do + metadata = described_class.observation_metadata( + name: "quality", run_name: "nightly", + item_metadata: { shared: "item", item_only: true }, + experiment_metadata: { shared: "run", run_only: true }, + dataset_id: "dataset-id", dataset_item_id: "item-id" + ) + + expect(metadata).to include( + shared: "run", item_only: true, run_only: true, + experiment_name: "quality", experiment_run_name: "nightly", + dataset_id: "dataset-id", dataset_item_id: "item-id" + ) + end + end +end diff --git a/spec/langfuse/experiment_result_spec.rb b/spec/langfuse/experiment_result_spec.rb index 406ebb7..dd183ca 100644 --- a/spec/langfuse/experiment_result_spec.rb +++ b/spec/langfuse/experiment_result_spec.rb @@ -41,6 +41,7 @@ expect(result.run_evaluations).to eq([]) expect(result.run_name).to be_nil expect(result.description).to be_nil + expect(result.experiment_id).to be_nil expect(result.dataset_run_id).to be_nil expect(result.dataset_run_url).to be_nil end @@ -57,8 +58,10 @@ it "stores dataset_run_id and dataset_run_url" do result = described_class.new( name: "exp", item_results: [], + experiment_id: "experiment-123", dataset_run_id: "run-123", dataset_run_url: "https://example.com/run/123" ) + expect(result.experiment_id).to eq("experiment-123") expect(result.dataset_run_id).to eq("run-123") expect(result.dataset_run_url).to eq("https://example.com/run/123") end diff --git a/spec/langfuse/experiment_runner_spec.rb b/spec/langfuse/experiment_runner_spec.rb index 33001ab..2d5ba6e 100644 --- a/spec/langfuse/experiment_runner_spec.rb +++ b/spec/langfuse/experiment_runner_spec.rb @@ -46,6 +46,7 @@ expect(result).to be_a(Langfuse::ExperimentResult) expect(result.name).to eq("my-exp") + expect(result.experiment_id).to match(/\A[0-9a-f]{16}\z/) end it "generates run_name from name and timestamp by default" do @@ -285,7 +286,7 @@ expect(result.dataset_run_url).to be_nil end - it "uses dataset_id from the first successfully linked item" do + it "uses dataset_id from the link that establishes the run identity" do items = [ Langfuse::DatasetItemClient.new( { "id" => "item-1", "datasetId" => "ds-1", @@ -299,13 +300,8 @@ ) ] - call_count = 0 - allow(mock_client).to receive(:create_dataset_run_item) do - call_count += 1 - raise StandardError, "link error" if call_count == 1 - - { "datasetRunId" => "run-1" } - end + allow(mock_client).to receive(:create_dataset_run_item) + .and_return({ "datasetRunId" => "run-1" }) allow(mock_client).to receive(:dataset_run_url) .with(dataset_id: "ds-1", dataset_run_id: "run-1") .and_return("https://example.com/datasets/ds-1/runs/run-1") @@ -505,7 +501,11 @@ runner.execute expect(mock_client).to have_received(:create_score).with( - hash_including(name: "score", value: 0.9) + hash_including( + name: "score", value: 0.9, + trace_id: a_string_matching(/\A[0-9a-f]{32}\z/), + observation_id: a_string_matching(/\A[0-9a-f]{16}\z/) + ) ) end @@ -843,6 +843,10 @@ end context "when dataset linking fails" do + before do + allow(logger).to receive(:warn).and_return(true) + end + let(:dataset_item) do Langfuse::DatasetItemClient.new( { "id" => "item-1", "datasetId" => "ds-1", @@ -852,17 +856,195 @@ end it "logs warning and continues" do + root_attributes = nil allow(mock_client).to receive(:create_dataset_run_item) .and_raise(StandardError, "link error") runner = described_class.new( - client: mock_client, name: "test", items: [dataset_item], task: ->(_) { "a" } + client: mock_client, name: "test", items: [dataset_item], + task: lambda { |_item| + root_attributes = OpenTelemetry::Trace.current_span.attributes.dup + "a" + } ) result = runner.execute expect(result.item_results.first.success?).to be true + expect(result.experiment_id).to match(/\A[0-9a-f]{16}\z/) + expect(root_attributes["langfuse.experiment.id"]).to eq(result.experiment_id) + expect(logger).to have_received(:warn).with(/Dataset run item linking failed/) + end + + it "reuses a known dataset run ID after a later link failure" do + second_item = Langfuse::DatasetItemClient.new( + { "id" => "item-2", "datasetId" => "ds-1", + "input" => { "q" => "second" }, "expectedOutput" => "a2" }, + client: mock_client + ) + call_count = 0 + allow(mock_client).to receive(:create_dataset_run_item) do + call_count += 1 + raise StandardError, "link error" if call_count == 2 + + { "datasetRunId" => "known-run-id" } + end + experiment_ids = [] + runner = described_class.new( + client: mock_client, name: "test", items: [dataset_item, second_item], + task: lambda { |_item| + experiment_ids << OpenTelemetry::Trace.current_span.attributes["langfuse.experiment.id"] + "a" + } + ) + + result = runner.execute + + expect(experiment_ids).to eq(%w[known-run-id known-run-id]) + expect(result.experiment_id).to eq("known-run-id") + expect(logger).to have_received(:warn).with(/Dataset run item linking failed/) + end + + it "keeps the fallback experiment ID after a later link succeeds" do + second_item = Langfuse::DatasetItemClient.new( + { "id" => "item-2", "datasetId" => "ds-1", + "input" => { "q" => "second" }, "expectedOutput" => "a2" }, + client: mock_client + ) + call_count = 0 + allow(mock_client).to receive(:create_dataset_run_item) do + call_count += 1 + raise StandardError, "link error" if call_count == 1 + + { "datasetRunId" => "later-run-id" } + end + experiment_ids = [] + runner = described_class.new( + client: mock_client, name: "test", items: [dataset_item, second_item], + task: lambda { |_item| + experiment_ids << OpenTelemetry::Trace.current_span.attributes["langfuse.experiment.id"] + "a" + } + ) + + result = runner.execute + + expect(experiment_ids).to all(eq(result.experiment_id)) + expect(result.experiment_id).to match(/\A[0-9a-f]{16}\z/) + expect(result.dataset_run_id).to be_nil + expect(result.dataset_run_url).to be_nil + expect(mock_client).to have_received(:create_dataset_run_item).twice + expect(mock_client).not_to have_received(:dataset_run_url) expect(logger).to have_received(:warn).with(/Dataset run item linking failed/) end end + + context "with v4 experiment attributes" do + it "sets shared identity on local item roots and child observations" do # rubocop:disable RSpec/ExampleLength + root_attributes = nil + child_attributes = nil + items = [{ input: { question: "What?" }, expected_output: false, + metadata: { difficulty: "easy" } }] + task = lambda do |_item| + root_attributes = OpenTelemetry::Trace.current_span.attributes.dup + Langfuse.observe("experiment-child") do |child| + child_attributes = child.otel_span.attributes.dup + end + "answer" + end + + result = described_class.new( + client: mock_client, name: "quality", items: items, task: task, + run_name: "nightly", description: "quality check", metadata: { model: "test" } + ).execute + + expected_item_id = Langfuse::ExperimentAttributes.generate_item_id(items.first[:input]) + expected_shared = { + "langfuse.experiment.id" => result.experiment_id, + "langfuse.experiment.name" => "nightly", + "langfuse.experiment.item.id" => expected_item_id, + "langfuse.experiment.item.root_observation_id" => result.item_results.first.observation_id, + "langfuse.experiment.metadata.model" => "test", + "langfuse.experiment.item.metadata.difficulty" => "easy", + "langfuse.environment" => "sdk-experiment" + } + expect(root_attributes).to include(expected_shared) + expect(child_attributes).to include(expected_shared) + expect(root_attributes).to include( + "langfuse.experiment.description" => "quality check", + "langfuse.experiment.item.expected_output" => "false" + ) + expect(child_attributes).not_to include("langfuse.experiment.description") + expect(child_attributes).not_to include("langfuse.experiment.item.expected_output") + end + + it "uses the managed dataset run and item identifiers" do # rubocop:disable RSpec/ExampleLength + root_attributes = nil + child_attributes = nil + dataset_item = Langfuse::DatasetItemClient.new( + { "id" => "item-1", "datasetId" => "dataset-1", "input" => "question", + "expectedOutput" => "answer", "metadata" => { "shared" => "item", "item_only" => true } }, + client: mock_client + ) + allow(mock_client).to receive(:create_dataset_run_item) + .and_return({ "datasetRunId" => "dataset-run-1" }) + + result = described_class.new( + client: mock_client, name: "quality", items: [dataset_item], run_name: "nightly", + metadata: { shared: "run", run_only: true }, + task: lambda { |_item| + root_attributes = OpenTelemetry::Trace.current_span.attributes.dup + Langfuse.observe("experiment-child") { |child| child_attributes = child.otel_span.attributes.dup } + "answer" + } + ).execute + + expected_shared = { + "langfuse.experiment.id" => "dataset-run-1", + "langfuse.experiment.dataset.id" => "dataset-1", + "langfuse.experiment.item.id" => "item-1", + "langfuse.experiment.item.root_observation_id" => result.item_results.first.observation_id + } + expect(result.experiment_id).to eq("dataset-run-1") + expect(root_attributes).to include(expected_shared) + expect(child_attributes).to include(expected_shared) + expect(root_attributes).to include( + "langfuse.observation.metadata.shared" => "run", + "langfuse.observation.metadata.item_only" => "true", + "langfuse.observation.metadata.run_only" => "true", + "langfuse.observation.metadata.experiment_name" => "quality", + "langfuse.observation.metadata.experiment_run_name" => "nightly", + "langfuse.observation.metadata.dataset_id" => "dataset-1", + "langfuse.observation.metadata.dataset_item_id" => "item-1" + ) + end + + it "does not leak experiment context between concurrent runs" do + ready = Queue.new + release = Queue.new + captured = Queue.new + threads = %w[first second].map do |run_name| + Thread.new do + runner = described_class.new( + client: mock_client, name: "quality", items: [{ input: run_name }], run_name: run_name, + task: lambda { |_item| + ready << true + release.pop + Langfuse.observe("child") do |child| + captured << [run_name, child.otel_span.attributes["langfuse.experiment.name"]] + end + "answer" + } + ) + runner.execute + end + end + + 2.times { ready.pop } + 2.times { release << true } + threads.each(&:join) + + expect(2.times.map { captured.pop }).to contain_exactly(%w[first first], %w[second second]) + end + end end end diff --git a/spec/langfuse/otel_attributes_spec.rb b/spec/langfuse/otel_attributes_spec.rb index 11d27d8..6a81a79 100644 --- a/spec/langfuse/otel_attributes_spec.rb +++ b/spec/langfuse/otel_attributes_spec.rb @@ -37,6 +37,20 @@ expect(described_class::ENVIRONMENT).to eq("langfuse.environment") expect(described_class::IS_APP_ROOT).to eq("langfuse.internal.is_app_root") end + + it "defines v4 experiment attribute constants" do + expect(described_class::EXPERIMENT_ID).to eq("langfuse.experiment.id") + expect(described_class::EXPERIMENT_NAME).to eq("langfuse.experiment.name") + expect(described_class::EXPERIMENT_DESCRIPTION).to eq("langfuse.experiment.description") + expect(described_class::EXPERIMENT_METADATA).to eq("langfuse.experiment.metadata") + expect(described_class::EXPERIMENT_DATASET_ID).to eq("langfuse.experiment.dataset.id") + expect(described_class::EXPERIMENT_ITEM_ID).to eq("langfuse.experiment.item.id") + expect(described_class::EXPERIMENT_ITEM_EXPECTED_OUTPUT) + .to eq("langfuse.experiment.item.expected_output") + expect(described_class::EXPERIMENT_ITEM_METADATA).to eq("langfuse.experiment.item.metadata") + expect(described_class::EXPERIMENT_ITEM_ROOT_OBSERVATION_ID) + .to eq("langfuse.experiment.item.root_observation_id") + end end describe ".serialize" do diff --git a/spec/langfuse/propagation_spec.rb b/spec/langfuse/propagation_spec.rb index 9e526e9..2758aca 100644 --- a/spec/langfuse/propagation_spec.rb +++ b/spec/langfuse/propagation_spec.rb @@ -10,6 +10,29 @@ end end + describe "._with_experiment_attributes" do + it "sets current and child spans without leaking into later siblings" do + attributes = { "langfuse.experiment.id" => "experiment-1" } + + Langfuse.observe("root") do |root| + described_class._with_experiment_attributes(attributes) do + child = root.start_observation("inside") + expect(root.otel_span.attributes).to include(attributes) + expect(child.otel_span.attributes).to include(attributes) + child.end + end + + sibling = root.start_observation("outside") + expect(sibling.otel_span.attributes).not_to include("langfuse.experiment.id") + sibling.end + end + end + + it "executes directly when attributes are empty" do + expect(described_class._with_experiment_attributes({}) { "result" }).to eq("result") + end + end + shared_context "with baggage mock" do before do baggage_module = Module.new do diff --git a/spec/langfuse/traced_execution_spec.rb b/spec/langfuse/traced_execution_spec.rb index 12d8c95..979839d 100644 --- a/spec/langfuse/traced_execution_spec.rb +++ b/spec/langfuse/traced_execution_spec.rb @@ -111,47 +111,52 @@ expect(error).to be_nil end - it "yields span and trace_id to the pre-task hook" do + it "passes span and trace_id to the context preparer" do yielded_span = nil yielded_trace_id = nil described_class.call( trace_name: "test-trace", input: {}, + prepare_context: lambda { |span, trace_id| + yielded_span = span + yielded_trace_id = trace_id + {} + }, task: ->(_span) { "result" } - ) do |span, trace_id| - yielded_span = span - yielded_trace_id = trace_id - end + ) expect(yielded_span).to be_a(Langfuse::BaseObservation) expect(yielded_trace_id).to be_a(String) expect(yielded_trace_id.length).to eq(32) end - it "executes the pre-task hook before the task" do + it "executes the context preparer before the task" do call_order = [] described_class.call( trace_name: "test-trace", input: {}, + prepare_context: lambda { |_span, _trace_id| + call_order << :context + {} + }, task: lambda { |_span| call_order << :task "result" } - ) do |_span, _trace_id| - call_order << :hook - end + ) - expect(call_order).to eq(%i[hook task]) + expect(call_order).to eq(%i[context task]) end - it "still captures task error when pre-task hook is given" do + it "still captures task error when a context preparer is given" do _output, _trace_id, _observation_id, error = described_class.call( trace_name: "test-trace", input: {}, + prepare_context: ->(_span, _trace_id) { {} }, task: ->(_span) { raise StandardError, "task failed" } - ) { |_span, _trace_id| } + ) expect(error).to be_a(StandardError) expect(error.message).to eq("task failed")