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
35 changes: 32 additions & 3 deletions lib/braintrust/eval/runner.rb
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,11 @@ def initialize(eval_context)
@eval_context = eval_context
@tracer = eval_context.tracer_provider.tracer("braintrust-eval")

# Whether any scorer/classifier can receive `trace:`. Computed once: the
# per-case flush that makes traces queryable over BTQL is a real network
# wait, so it's only worth paying when something will actually consume it.
@needs_trace = (eval_context.scorers + eval_context.classifiers).any? { |c| wants_trace?(c) }

# Mutexes for thread-safe result collection
@score_mutex = Mutex.new
@classification_mutex = Mutex.new
Expand Down Expand Up @@ -55,6 +60,11 @@ def run(parallelism: 1)
# Convert Queue to Array after all threads complete
error_array = [].tap { |a| a << errors.pop until errors.empty? }

# Deliver any spans still buffered. Callers that supply their own
# tracer_provider get no at_exit hook (see Trace.setup), so without this
# the tail of the run would wait on the processor's schedule delay.
flush_spans

# Calculate duration
duration = Time.now - start_time

Expand Down Expand Up @@ -108,9 +118,11 @@ def run_eval_case(kase, errors)
next
end

# Flush spans so they're queryable via BTQL, then build trace
eval_context.tracer_provider.force_flush if eval_context.tracer_provider.respond_to?(:force_flush)
kase.trace = build_trace(eval_span)
# Build the trace, then flush spans so they're queryable via BTQL.
# Both are skipped unless a scorer/classifier declared `trace:` and the
# trace is actually resolvable (build_trace returns nil in local-only mode).
kase.trace = build_trace(eval_span) if @needs_trace
flush_spans if kase.trace

# Run scorers
begin
Expand Down Expand Up @@ -229,6 +241,23 @@ def run_scorer(scorer, scorer_kwargs, scorer_input)
end
end

# Whether a callable can receive `trace:`: it either declares the keyword
# or accepts arbitrary kwargs and may forward it. Reuses the #call_parameters
# introspection Internal::Callable::KeywordFilter already depends on.
# @param callable [Scorer, Classifier]
# @return [Boolean]
def wants_trace?(callable)
return true unless callable.respond_to?(:call_parameters)

callable.call_parameters.any? { |type, name| name == :trace || type == :keyrest }
end

# Force the tracer provider to export buffered spans, if it supports it.
def flush_spans
provider = eval_context.tracer_provider
provider.force_flush if provider.respond_to?(:force_flush)
end

# Build a lazy Trace for a case, backed by BTQL.
# Returns nil when state or experiment_id are unavailable (local-only mode).
# @param eval_span [OpenTelemetry::Trace::Span] The eval span for this case
Expand Down
4 changes: 2 additions & 2 deletions lib/braintrust/scorer.rb
Original file line number Diff line number Diff line change
Expand Up @@ -152,10 +152,10 @@ def wrap_block(block)
case block.arity
when 3
Log.warn_once(:scorer_positional_3, "Scorer with positional params (input, expected, output) is deprecated. Use keyword args: |input:, expected:, output:| instead.")
->(**kw) { block.call(kw[:input], kw[:expected], kw[:output]) }
->(input: nil, expected: nil, output: nil) { block.call(input, expected, output) }
when 4, -4, -1
Log.warn_once(:scorer_positional_4, "Scorer with positional params (input, expected, output, metadata) is deprecated. Use keyword args: |input:, expected:, output:, metadata:| instead.")
->(**kw) { block.call(kw[:input], kw[:expected], kw[:output], kw[:metadata]) }
->(input: nil, expected: nil, output: nil, metadata: nil) { block.call(input, expected, output, metadata) }
else
raise ArgumentError, "Scorer must accept keyword args or 3-4 positional params (got arity #{block.arity})"
end
Expand Down
141 changes: 141 additions & 0 deletions test/braintrust/eval/runner_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -1301,6 +1301,117 @@ def test_trace_works_with_parallelism
end
end

# ============================================
# Runner#run tests - per-case flush gating
# ============================================
#
# The per-case force_flush exists only so a `trace:`-consuming scorer/classifier
# can BTQL-query that case's spans. It is a real network wait (10-25s+ against a
# live backend), so it must not run when nothing will consume the trace.
#
# Counts below include the single end-of-run flush, which always happens.

def test_no_per_case_flush_when_no_scorer_declares_trace
rig = setup_otel_test_rig
flushes = count_force_flushes(rig.tracer_provider)

scorer = Braintrust::Scorer.new("simple") { |output:, expected:| (output == expected) ? 1.0 : 0.0 }
result = run_flush_eval(rig, scorers: [scorer], cases: 3)

assert result.success?
assert_equal 1, flushes.length, "expected only the end-of-run flush"
end

def test_no_per_case_flush_for_legacy_positional_scorer
rig = setup_otel_test_rig
flushes = count_force_flushes(rig.tracer_provider)

# Deprecated positional signature is wrapped internally; the wrapper must not
# look like **kwargs, or it would be treated as a possible trace consumer.
scorer = Braintrust::Scorer.new("legacy") { |_input, expected, output| (output == expected) ? 1.0 : 0.0 }
result = run_flush_eval(rig, scorers: [scorer], cases: 3)

assert result.success?
assert_equal 1, flushes.length, "expected only the end-of-run flush"
end

def test_per_case_flush_when_scorer_declares_trace
rig = setup_otel_test_rig
flushes = count_force_flushes(rig.tracer_provider)

scorer = Braintrust::Scorer.new("trace_reader") { |output:, trace:| trace.nil? ? 0.0 : 1.0 }
result = run_flush_eval(rig, scorers: [scorer], cases: 3)

assert result.success?
assert_equal 4, flushes.length, "expected one flush per case plus the end-of-run flush"
end

def test_per_case_flush_when_only_a_classifier_declares_trace
rig = setup_otel_test_rig
flushes = count_force_flushes(rig.tracer_provider)

scorer = Braintrust::Scorer.new("simple") { |output:| 1.0 }
classifier = Braintrust::Classifier.new("kind") { |output:, trace:| {name: "kind", id: "seen", label: "Seen"} }
result = run_flush_eval(rig, scorers: [scorer], classifiers: [classifier], cases: 2)

assert result.success?
assert_equal 3, flushes.length, "expected one flush per case plus the end-of-run flush"
end

def test_per_case_flush_when_scorer_accepts_arbitrary_kwargs
rig = setup_otel_test_rig
flushes = count_force_flushes(rig.tracer_provider)

# **kwargs might forward trace:, so stay conservative and keep flushing.
scorer = Braintrust::Scorer.new("splat") { |**kwargs| 1.0 }
result = run_flush_eval(rig, scorers: [scorer], cases: 2)

assert result.success?
assert_equal 3, flushes.length, "expected one flush per case plus the end-of-run flush"
end

def test_no_per_case_flush_when_trace_is_unavailable
rig = setup_otel_test_rig
flushes = count_force_flushes(rig.tracer_provider)
received_trace = :not_called

scorer = Braintrust::Scorer.new("trace_reader") { |output:, trace:|
received_trace = trace
1.0
}

# No experiment_id: build_trace returns nil, so the flush buys nothing.
context = Braintrust::Eval::Context.build(
task: ->(input:) { input.upcase },
scorers: [scorer],
cases: [{input: "a"}, {input: "b"}],
state: rig.state,
tracer_provider: rig.tracer_provider
)
result = Braintrust::Eval::Runner.new(context).run

assert result.success?
assert_nil received_trace
assert_equal 1, flushes.length, "expected only the end-of-run flush"
end

def test_run_flushes_once_after_all_cases
rig = setup_otel_test_rig
flushed_after = nil
scored = 0

scorer = Braintrust::Scorer.new("counter") { |output:|
scored += 1
1.0
}
rig.tracer_provider.define_singleton_method(:force_flush) { |timeout: nil| flushed_after = scored }

result = run_flush_eval(rig, scorers: [scorer], cases: 3)

assert result.success?
assert_equal 3, flushed_after, "end-of-run flush should happen after every case has been scored"
end

# ============================================
# Runner#run tests - structured scorer returns
# ============================================
Expand Down Expand Up @@ -2064,6 +2175,36 @@ def test_runner_parameters_with_parallelism
assert_equal 3, params.length
assert(params.all? { |p| p == {"model" => "gpt-4"} })
end

private

# Replace the provider's #force_flush with a recorder. Returns the array of
# recorded calls so a test can assert how many flushes a run performed.
def count_force_flushes(tracer_provider)
calls = []
tracer_provider.define_singleton_method(:force_flush) do |timeout: nil|
calls << timeout
OpenTelemetry::SDK::Trace::Export::SUCCESS
end
calls
end

# Run an eval with a fully-populated experiment context, so build_trace resolves.
def run_flush_eval(rig, scorers:, classifiers: [], cases: 1)
context = Braintrust::Eval::Context.build(
task: ->(input:) { input.upcase },
scorers: scorers,
classifiers: classifiers,
cases: (1..cases).map { |i| {input: "case-#{i}", expected: "CASE-#{i}"} },
experiment_id: "exp-123",
experiment_name: "test-experiment",
project_id: "proj-456",
project_name: "test-project",
state: rig.state,
tracer_provider: rig.tracer_provider
)
Braintrust::Eval::Runner.new(context).run
end
end

class Braintrust::Eval::RunnerClassifierTest < Minitest::Test
Expand Down
32 changes: 32 additions & 0 deletions test/braintrust/scorer_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,38 @@ def test_scorer_return_array
assert_equal({name: "metric2", score: 0.8}, result[1])
end

# ============================================
# call_parameters introspection
# ============================================
#
# Callers introspect #call_parameters to decide what a scorer can consume.
# Eval::Runner uses it to skip its per-case span flush when nothing declares
# `trace:`, so the legacy positional wrappers must report the keywords they
# actually take rather than an opaque **kwargs.

def test_legacy_3_param_block_reports_declared_keywords
suppress_logs do
scorer = Braintrust::Scorer.new("legacy3") { |input, expected, output| 1.0 }

assert_equal [[:key, :input], [:key, :expected], [:key, :output]], scorer.call_parameters
end
end

def test_legacy_4_param_block_reports_declared_keywords
suppress_logs do
scorer = Braintrust::Scorer.new("legacy4") { |input, expected, output, metadata| 1.0 }

assert_equal [[:key, :input], [:key, :expected], [:key, :output], [:key, :metadata]],
scorer.call_parameters
end
end

def test_keyword_block_reports_its_own_parameters
scorer = Braintrust::Scorer.new("kw") { |output:, trace:| 1.0 }

assert_equal [[:keyreq, :output], [:keyreq, :trace]], scorer.call_parameters
end

# ============================================
# Validation
# ============================================
Expand Down
Loading