diff --git a/CHANGELOG.md b/CHANGELOG.md index 4a5597d..5b51de3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,34 @@ # Changelog -## [Unreleased] +## [0.3.0] - 2026-09-12 + +### Added + +- `RubyLLM.with_agent` takes `attributes:`, `on_trace:` and `synchronous:`. + `attributes` are merged onto the root span of every trace recorded in the + block, with the `agent.*` identity keys taking precedence, so an evaluation + can stamp its run and result identifiers on the traces it causes. `on_trace` + receives each trace the reporter accepted for delivery, so the caller can + keep the `trace_id` of that delivery attempt; a trace dropped by + `sample_rate` or a disabled configuration is never announced. Acceptance is + not ingestion: a delivery that fails afterwards is logged by the reporter, + not surfaced to the callback. `synchronous: true` delivers in + the calling thread for that scope only, through the same sampling and + configuration checks as ordinary delivery; the shared asynchronous + configuration is untouched, and a delivery failure follows the reporter's + existing logging policy rather than raising. A turn keeps the scope it + started under, so a turn left open by a pending tool call and closed later + by `flush!` still reports with that scope's agent, attributes and callback, + and a turn that started outside any scope never adopts a later one. + Nested scopes restore the previous context, including when the block + raises. A callback that raises is logged by exception class and the trace + is still delivered. (#5) +- `Reporter#report` takes `sync: true` to deliver in the calling thread for + that call only, keeping the enabled, configured and sampling checks that + `report_now` skips, and returns whether the traces were accepted. + `BatchingReporter#report` delivers that call's traces in the calling + thread when asked the same, leaving its buffer on its own schedule. The + RubyLLM adapter now requires core `~> 0.3` for it. ### Fixed diff --git a/adapters/ruby_llm/README.md b/adapters/ruby_llm/README.md index 311b22e..000dcac 100644 --- a/adapters/ruby_llm/README.md +++ b/adapters/ruby_llm/README.md @@ -92,3 +92,44 @@ an explicit `flush!`. ```bash bundle exec rake test ``` + +## Correlating evaluation traces + +Use a separate identity for judge calls and attach stable evaluation identifiers +without changing the application's default agent resolver: + +```ruby +trace_ids = [] +ActiveAgents::Telemetry::RubyLLM.with_agent( + "EvaluationJudge", action: "score", + attributes: { "eval.run_id" => run_id, "eval.result_id" => result_id }, + on_trace: ->(trace) { trace_ids << trace.trace_id }, + synchronous: true +) do + judge_chat.ask(prompt) +end +``` + +The callback receives each trace the reporter accepted for delivery, so an +evaluation result can retain the trace ID of that delivery attempt. Acceptance is +not ingestion: the trace passed the enabled, configured and sampling checks, but +a delivery that then fails is logged by the reporter rather than announced here, +and an asynchronous delivery can outlive a process that exits right away. A trace +dropped by `sample_rate` or by a disabled configuration is never announced. +`synchronous: true` delivers in the calling thread for this scope only, through +the same sampling and configuration checks as ordinary delivery; it does not +mutate the shared async configuration. Normal reporter error logging still +applies: synchronous delivery does not turn telemetry failures into application +exceptions. A turn keeps the scope it started under, so a turn left open by a +pending tool call and closed later by `flush!` still reports with this agent, +attributes and callback, and a turn that started outside any scope never adopts +one. Context +is restored after the block, including when it raises, and nested scopes use +their own attributes. A callback error is logged by exception class without +dropping the trace. Attribute redaction and content-capture settings continue to +apply. + +Report publication and trace ingestion are separate operations. Persist each run +and result ID in the evaluation report and attach the same IDs to its response and +judge trace attributes. The application chooses whether to publish full report +content; this API does not enable body capture or upload reports automatically. diff --git a/adapters/ruby_llm/activeagents-telemetry-ruby_llm.gemspec b/adapters/ruby_llm/activeagents-telemetry-ruby_llm.gemspec index 2ebd04e..4501f01 100644 --- a/adapters/ruby_llm/activeagents-telemetry-ruby_llm.gemspec +++ b/adapters/ruby_llm/activeagents-telemetry-ruby_llm.gemspec @@ -31,6 +31,6 @@ Gem::Specification.new do |spec| spec.files = Dir["lib/**/*.rb", "README.md"] spec.require_paths = [ "lib" ] - spec.add_dependency "activeagents-telemetry", "~> 0.1" + spec.add_dependency "activeagents-telemetry", "~> 0.3" spec.add_dependency "activesupport", ">= 7.0" end diff --git a/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm.rb b/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm.rb index f8687ca..ebe88b9 100644 --- a/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm.rb +++ b/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm.rb @@ -43,7 +43,7 @@ module RubyLLM DEFAULT_AGENT = { name: "RubyLLM::Chat", action: "chat" }.freeze - State = Struct.new(:depth, :started_at, :tool_spans, :rounds, :tokens, :chat_key) + State = Struct.new(:depth, :started_at, :tool_spans, :rounds, :tokens, :chat_key, :agent) class << self # Subscribes to RubyLLM's instrumentation. @@ -91,17 +91,30 @@ def reporter attr_writer :reporter - # Attributes traces inside the block to a named agent/action. - def with_agent(name, action: "chat") + # Attributes traces inside the block to a named agent/action. Correlation + # attributes and the callback apply to this scope only. Short-lived + # evaluation commands can deliver synchronously without changing the + # application's shared reporter configuration. + # + # A turn keeps the scope it started under, so a turn left open by a + # pending tool call and closed later by `flush!` still reports as this + # agent, and a turn that started outside any scope never adopts one. + # `on_trace` runs only for a trace the reporter accepted, which means + # it passed the enabled, configured and sampling checks: a delivery + # that then fails is logged by the reporter, not announced here. + def with_agent(name, action: "chat", attributes: {}, on_trace: nil, synchronous: false) previous = Thread.current[AGENT_KEY] - Thread.current[AGENT_KEY] = { name: name, action: action } + Thread.current[AGENT_KEY] = { + name: name, action: action, attributes: attributes.to_h.transform_keys(&:to_s), + on_trace: on_trace, synchronous: synchronous + } yield ensure Thread.current[AGENT_KEY] = previous end def state - Thread.current[STATE_KEY] ||= State.new(0, nil, [], 0, Span::ZERO_TOKENS.dup, nil) + Thread.current[STATE_KEY] ||= State.new(0, nil, [], 0, Span::ZERO_TOKENS.dup, nil, nil) end def clear_state @@ -126,7 +139,13 @@ def begin_round(payload) flush! if turn.rounds.positive? && (turn.chat_key != chat_key || turn_expired?(turn)) turn = state turn.chat_key = chat_key - turn.started_at ||= Time.now + if turn.started_at.nil? + turn.started_at = Time.now + # Captured once, on the turn's first round, whether or not a scope + # is active: a turn that started unscoped stays unscoped even when + # a later scope's chat is what flushes it. + turn.agent = Thread.current[AGENT_KEY] + end end turn.depth += 1 end @@ -161,7 +180,7 @@ def build_tool_span(payload, started_at, finished_at) private def report_turn(payload, turn) - agent = Thread.current[AGENT_KEY] || resolve_agent(payload) || DEFAULT_AGENT + agent = turn.agent || resolve_agent(payload) || DEFAULT_AGENT started_at = turn.started_at || Time.now finished_at = Time.now error = payload[:exception_object] @@ -172,12 +191,12 @@ def report_turn(payload, turn) resource_attributes: configuration.resource_attributes ) - root_attributes = { + root_attributes = (agent[:attributes] || {}).merge( "agent.class" => agent[:name], "agent.action" => agent[:action], "agent.provider" => payload[:provider].to_s, "agent.model" => payload[:model].to_s - } + ) root_attributes.merge!(conversation_attributes(payload)) if configuration.capture_bodies? root = trace.span( @@ -206,7 +225,14 @@ def report_turn(payload, turn) trace.add_span(tool_span) end - reporter.report(trace) + accepted = agent[:synchronous] ? reporter.report(trace, sync: true) : reporter.report(trace) + notify_trace(agent[:on_trace], trace) if accepted + end + + def notify_trace(callback, trace) + callback&.call(trace) + rescue StandardError => e + warn "[#{SDK_NAME}] on_trace failed: #{e.class}" end def turn_expired?(turn) diff --git a/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm/version.rb b/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm/version.rb index 220f1f1..28607ca 100644 --- a/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm/version.rb +++ b/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm/version.rb @@ -3,7 +3,7 @@ module ActiveAgents module Telemetry module RubyLLM - VERSION = "0.2.0" + VERSION = "0.3.0" end end end diff --git a/adapters/ruby_llm/test/test_evaluation_context.rb b/adapters/ruby_llm/test/test_evaluation_context.rb new file mode 100644 index 0000000..7f080e0 --- /dev/null +++ b/adapters/ruby_llm/test/test_evaluation_context.rb @@ -0,0 +1,120 @@ +# frozen_string_literal: true + +require "test_helper" + +class TestEvaluationContext < Minitest::Test + include RubyLLMTelemetryTestHelpers + + def test_correlates_the_actual_trace_and_separates_judge_identity + ids = [] + Adapter.with_agent("EvaluationJudge", action: "score", + attributes: { "eval.run_id" => "run-1", "eval.result_id" => "result-1", "agent.class" => "Wrong" }, + on_trace: ->(trace) { ids << trace.trace_id }) do + instrument("chat.ruby_llm", chat_payload) { nil } + end + + trace = traces.fetch(0) + root = spans_of(trace, "root").fetch(0) + assert_equal [ trace["trace_id"] ], ids + assert_equal "EvaluationJudge.score", root["name"] + assert_equal "EvaluationJudge", root["attributes"]["agent.class"] + assert_equal "run-1", root["attributes"]["eval.run_id"] + assert_equal "result-1", root["attributes"]["eval.result_id"] + end + + def test_context_is_restored_after_nested_judge_and_exception + Adapter.with_agent("Support", action: "respond", attributes: { "eval.run_id" => "outer" }) do + assert_raises(RuntimeError) do + Adapter.with_agent("Judge", action: "score", attributes: { "eval.run_id" => "inner" }) do + instrument("chat.ruby_llm", chat_payload) { nil } + raise "synthetic error" + end + end + instrument("chat.ruby_llm", chat_payload) { nil } + end + instrument("chat.ruby_llm", chat_payload) { nil } + + roots = traces.map { |trace| spans_of(trace, "root").first } + assert_equal [ "Judge.score", "Support.respond", "RubyLLM::Chat.chat" ], roots.map { |root| root["name"] } + assert_equal [ "inner", "outer", nil ], roots.map { |root| root["attributes"]["eval.run_id"] } + end + + def test_synchronous_scope_delivers_in_the_calling_thread_without_mutating_configuration + subscribe(async: true) + caller_thread = Thread.current + delivery_threads = [] + captured = posted + Adapter.reporter.define_singleton_method(:deliver) do |body| + delivery_threads << Thread.current + captured << body + end + + Adapter.with_agent("Support", synchronous: true) do + instrument("chat.ruby_llm", chat_payload) { nil } + assert_equal 1, posted.size + end + + assert_equal [ caller_thread ], delivery_threads + assert Adapter.configuration.async? + end + + def test_synchronous_scope_still_honours_sampling_and_skips_the_callback + configuration = ActiveAgents::Telemetry::Configuration.new + configuration.sample_rate = 0.0 + subscribe(configuration: configuration) + ids = [] + + Adapter.with_agent("Support", synchronous: true, on_trace: ->(trace) { ids << trace.trace_id }) do + instrument("chat.ruby_llm", chat_payload) { nil } + end + + assert_empty posted + assert_empty ids + end + + def test_an_unscoped_turn_never_adopts_the_scope_whose_chat_flushes_it + ids = [] + instrument("chat.ruby_llm", chat_payload(tool_call: true)) { nil } + assert_empty posted, "the unscoped turn stays open on its pending tool call" + + Adapter.with_agent("Judge", action: "score", attributes: { "eval.run_id" => "run-1" }, + on_trace: ->(trace) { ids << trace.trace_id }, synchronous: true) do + instrument("chat.ruby_llm", chat_payload(chat: Object.new)) { nil } + end + + roots = traces.map { |trace| spans_of(trace, "root").first } + assert_equal [ "RubyLLM::Chat.chat", "Judge.score" ], roots.map { |root| root["name"] } + assert_equal [ nil, "run-1" ], roots.map { |root| root["attributes"]["eval.run_id"] } + assert_equal [ traces.fetch(1)["trace_id"] ], ids, "only the judge's own trace reaches its callback" + end + + def test_a_turn_keeps_the_scope_it_started_under_when_flushed_later + ids = [] + pending = chat_payload(tool_call: true) + Adapter.with_agent("Judge", action: "score", attributes: { "eval.run_id" => "run-1" }, + on_trace: ->(trace) { ids << trace.trace_id }) do + instrument("chat.ruby_llm", pending) { nil } + end + assert_empty posted, "a turn with a pending tool call stays open" + + Adapter.flush!(pending) + + trace = traces.fetch(0) + root = spans_of(trace, "root").fetch(0) + assert_equal "Judge.score", root["name"] + assert_equal "run-1", root["attributes"]["eval.run_id"] + assert_equal [ trace["trace_id"] ], ids + end + + def test_callback_failure_does_not_discard_the_trace_or_expose_its_message + _, stderr = capture_io do + Adapter.with_agent("Support", on_trace: ->(_) { raise "private callback content" }) do + instrument("chat.ruby_llm", chat_payload) { nil } + end + end + + assert_equal 1, posted.size + assert_includes stderr, "on_trace failed: RuntimeError" + refute_includes stderr, "private callback content" + end +end diff --git a/docs/branches/evaluation-trace-context.md b/docs/branches/evaluation-trace-context.md new file mode 100644 index 0000000..dec1761 --- /dev/null +++ b/docs/branches/evaluation-trace-context.md @@ -0,0 +1,6 @@ +# Branch + +- Branch: `codex/evaluation-trace-context`. +- Base: `main` at `ccbab0ffd2b26722685d8d2f7f7ff74a45466373`. +- Scope: backward-compatible RubyLLM trace context and delivery control for evaluations. +- Release: consumers can pin this branch's reviewed commit before the next adapter release. diff --git a/docs/issues/evaluation-trace-context.md b/docs/issues/evaluation-trace-context.md new file mode 100644 index 0000000..38acc5b --- /dev/null +++ b/docs/issues/evaluation-trace-context.md @@ -0,0 +1,17 @@ +# Evaluation trace identity and completion + +Tool-less evaluation judge calls previously inherited an application's generic +tool-less identity. Applications also had no supported callback for retaining the +actual emitted trace ID in a scenario result, and short-lived commands inherited +asynchronous delivery that could outlive the process. + +`RubyLLM.with_agent` now accepts scope-local correlation attributes, an `on_trace` +callback and a synchronous delivery option. Existing calls remain compatible. +Nested/raising scopes restore their previous context; callback failures do not +drop traces or log callback message content. Body capture remains opt-in. + +Validation: core 39 tests / 100 assertions and RubyLLM adapter 20 tests / 71 +assertions pass on Ruby 4.0.2. New tests cover actual trace-ID correlation, judge +identity, nested context restoration, synchronous delivery under sampling, a +turn keeping the scope it started under, and callback failure. +Raw logs live in gitignored `tmp/`. diff --git a/docs/milestones/evaluation-trace-context.md b/docs/milestones/evaluation-trace-context.md new file mode 100644 index 0000000..d632483 --- /dev/null +++ b/docs/milestones/evaluation-trace-context.md @@ -0,0 +1,8 @@ +# Milestone: evaluation correlation + +- [x] Separate candidate and judge trace identity. +- [x] Expose actual completed trace IDs for report correlation. +- [x] Support blocking delivery attempts for short-lived commands. +- [x] Test restoration, callback failures, and existing adapter compatibility. +- [x] Bump both gems to 0.3.0 with a changelog entry, so the merge is releasable. +- [ ] Merge the PR, then tag `v0.3.0` to publish both gems through the release workflow. diff --git a/docs/pull-requests/evaluation-trace-context.md b/docs/pull-requests/evaluation-trace-context.md new file mode 100644 index 0000000..b64dce5 --- /dev/null +++ b/docs/pull-requests/evaluation-trace-context.md @@ -0,0 +1,12 @@ +# Draft PR: evaluation trace context + +Branch `codex/evaluation-trace-context` proposes scoped correlation attributes, +completed trace callbacks, and blocking delivery attempts for RubyLLM calls. +It keeps existing agent resolvers and ordinary asynchronous application tracing +compatible. Consumers can distinguish evaluation judges from application agents +and link report results to their actual response and judge traces. + +Core and adapter suites pass: 59 tests, 171 assertions, zero failures. Changed +Ruby files pass the Omakase lint configuration. No provider or collector network +requests are made by the tests. Versions and the changelog are bumped to 0.3.0 +on the branch; tagging `v0.3.0` after the merge publishes both gems. diff --git a/lib/activeagents/telemetry/batching_reporter.rb b/lib/activeagents/telemetry/batching_reporter.rb index 530bd6c..757c010 100644 --- a/lib/activeagents/telemetry/batching_reporter.rb +++ b/lib/activeagents/telemetry/batching_reporter.rb @@ -24,13 +24,28 @@ def initialize(configuration, **options) @shutdown = false end - # Enqueues a trace, flushing if the batch is full. - def report(traces) - return if @shutdown + # Enqueues a trace, flushing if the batch is full. `sync: true` skips the + # buffer: that call's traces are delivered in the calling thread before + # it returns, so a short-lived process can hand a trace over before it + # exits, and whatever the buffer already holds stays on its own schedule. + # @return [Boolean] whether any of the traces were accepted + def report(traces, sync: false) + return false if @shutdown accepted = normalize(traces).select { sample_trace? } - return if accepted.empty? - return unless configuration.enabled? && configuration.configured? + return false if accepted.empty? + return false unless configuration.enabled? && configuration.configured? + + if sync + body = begin + payload_for(accepted) + rescue StandardError => e + log("failed to build trace payload: #{e.class}: #{e.message}") + return false + end + deliver(body) + return true + end batch = nil @mutex.synchronize do @@ -39,7 +54,7 @@ def report(traces) start_flusher end deliver_batch(batch) if batch - nil + true end # Delivers everything buffered, blocking until done. diff --git a/lib/activeagents/telemetry/reporter.rb b/lib/activeagents/telemetry/reporter.rb index 45c77d3..9164a03 100644 --- a/lib/activeagents/telemetry/reporter.rb +++ b/lib/activeagents/telemetry/reporter.rb @@ -35,19 +35,29 @@ def initialize(configuration, sdk_name: SDK_NAME, sdk_version: VERSION, sample: # @param traces [Trace, Hash, Array] traces to deliver — # Trace objects or already-serialized trace hashes - # @return [void] - def report(traces) + # @param sync [Boolean] deliver in the calling thread for this call + # only; every other call keeps the configured async behaviour. The + # enabled, configured and sampling checks still apply, unlike + # #report_now. + # @return [Boolean] whether the traces were accepted for delivery. A + # delivery failure is logged rather than surfaced here, so true means + # the traces passed the enabled, configured and sampling checks. + def report(traces, sync: false) traces = normalize(traces) - return if traces.empty? - return unless configuration.enabled? && configuration.configured? - return unless sample_trace? + return false if traces.empty? + return false unless configuration.enabled? && configuration.configured? + return false unless sample_trace? body = payload_for(traces) - configuration.async? ? Thread.new { deliver(body) } : deliver(body) - nil + if sync || !configuration.async? + deliver(body) + else + Thread.new { deliver(body) } + end + true rescue StandardError => e log("failed to build trace payload: #{e.class}: #{e.message}") - nil + false end # Blocking delivery, for tests and for at-exit flushes. diff --git a/lib/activeagents/telemetry/version.rb b/lib/activeagents/telemetry/version.rb index b5747b4..cf6368d 100644 --- a/lib/activeagents/telemetry/version.rb +++ b/lib/activeagents/telemetry/version.rb @@ -2,6 +2,6 @@ module ActiveAgents module Telemetry - VERSION = "0.2.0" + VERSION = "0.3.0" end end diff --git a/test/test_parity.rb b/test/test_parity.rb index 1b7cdf8..c912d3a 100644 --- a/test/test_parity.rb +++ b/test/test_parity.rb @@ -81,7 +81,7 @@ def test_a_local_store_counts_as_configured_without_endpoint_or_key config = ActiveAgents::Telemetry::Configuration.new refute config.configured? - config.local_store = ->(_trace, _sdk) {} + config.local_store = ->(_trace, _sdk) { } assert config.configured? end end @@ -165,6 +165,39 @@ def test_buffers_until_batch_size_then_delivers_together assert_equal 3, captured.first["traces"].size end + def test_sync_delivers_its_own_traces_in_the_calling_thread_and_leaves_the_buffer_alone + reporter = ActiveAgents::Telemetry::BatchingReporter.new(fresh_configuration(async: true)) + deliveries = [] + reporter.define_singleton_method(:deliver) { |body| deliveries << [ Thread.current, body["traces"].size ] } + + reporter.report(build_trace) + assert_empty deliveries, "an ordinary trace stays buffered" + + assert_equal true, reporter.report(build_trace, sync: true) + assert_equal [ [ Thread.current, 1 ] ], deliveries, "only the synchronous call's trace went out, in this thread" + + reporter.flush + assert_equal 2, deliveries.size, "the buffered trace was left for the next flush" + assert_equal 1, deliveries.last[1] + reporter.shutdown + end + + def test_sync_reports_a_trace_it_could_not_serialize_as_not_accepted + reporter, captured = batching_reporter(fresh_configuration(logger: Logger.new(File::NULL))) + trace = build_trace + trace.define_singleton_method(:to_h) { raise "unserializable" } + + assert_equal false, reporter.report(trace, sync: true) + assert_empty captured + end + + def test_sync_still_honours_sampling + reporter, captured = batching_reporter(fresh_configuration(sample_rate: 0.0)) + + assert_equal false, reporter.report(build_trace, sync: true) + assert_empty captured + end + def test_flush_delivers_a_partial_batch reporter, captured = batching_reporter reporter.report(build_trace) diff --git a/test/test_telemetry.rb b/test/test_telemetry.rb index 855934a..3c89639 100644 --- a/test/test_telemetry.rb +++ b/test/test_telemetry.rb @@ -151,8 +151,24 @@ def test_symbol_keyed_trace_hashes_are_stringified def test_sampling_drops_traces reporter, captured = capturing_reporter(fresh_configuration(sample_rate: 0.0)) - reporter.report(build_trace) + assert_equal false, reporter.report(build_trace) + assert_empty captured + end + + def test_sync_delivers_in_the_calling_thread_and_reports_acceptance + reporter = ActiveAgents::Telemetry::Reporter.new(fresh_configuration(async: true)) + threads = [] + reporter.define_singleton_method(:deliver) { |_body| threads << Thread.current } + + assert_equal true, reporter.report(build_trace, sync: true) + assert_equal [ Thread.current ], threads + end + + def test_sync_still_honours_sampling + reporter, captured = capturing_reporter(fresh_configuration(sample_rate: 0.0)) + + assert_equal false, reporter.report(build_trace, sync: true) assert_empty captured end @@ -162,6 +178,7 @@ def test_a_broken_endpoint_never_raises reporter = ActiveAgents::Telemetry::Reporter.new(config) config.logger = Logger.new(File::NULL) - assert_nil reporter.report(build_trace) + # The failure is logged inside delivery; the trace still counts as accepted. + assert_equal true, reporter.report(build_trace) end end