From dd411068bd9998d4e9ecdaf6d3245e29ee4ec320 Mon Sep 17 00:00:00 2001 From: Justin Bowen Date: Wed, 9 Sep 2026 13:28:07 -0700 Subject: [PATCH 1/6] feat(ruby_llm): correlate evaluation traces and control delivery --- adapters/ruby_llm/README.md | 31 ++++++++ .../lib/activeagents/telemetry/ruby_llm.rb | 25 +++++-- .../ruby_llm/test/test_evaluation_context.rb | 72 +++++++++++++++++++ docs/branches/evaluation-trace-context.md | 6 ++ docs/issues/evaluation-trace-context.md | 16 +++++ docs/milestones/evaluation-trace-context.md | 7 ++ .../pull-requests/evaluation-trace-context.md | 11 +++ 7 files changed, 162 insertions(+), 6 deletions(-) create mode 100644 adapters/ruby_llm/test/test_evaluation_context.rb create mode 100644 docs/branches/evaluation-trace-context.md create mode 100644 docs/issues/evaluation-trace-context.md create mode 100644 docs/milestones/evaluation-trace-context.md create mode 100644 docs/pull-requests/evaluation-trace-context.md diff --git a/adapters/ruby_llm/README.md b/adapters/ruby_llm/README.md index 311b22e..aa01a3f 100644 --- a/adapters/ruby_llm/README.md +++ b/adapters/ruby_llm/README.md @@ -92,3 +92,34 @@ 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 the completed trace before delivery, so an evaluation +result can retain the exact trace ID. `synchronous: true` waits for the existing +reporter's delivery attempt in this scope; it does not mutate the shared async +configuration. Normal reporter error logging still applies: synchronous delivery +does not turn telemetry failures into application exceptions. 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/lib/activeagents/telemetry/ruby_llm.rb b/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm.rb index f8687ca..9e4c924 100644 --- a/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm.rb +++ b/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm.rb @@ -91,10 +91,16 @@ 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. + 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 @@ -172,12 +178,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 +212,14 @@ def report_turn(payload, turn) trace.add_span(tool_span) end - reporter.report(trace) + notify_trace(agent[:on_trace], trace) + agent[:synchronous] ? reporter.report_now(trace) : reporter.report(trace) + 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/test/test_evaluation_context.rb b/adapters/ruby_llm/test/test_evaluation_context.rb new file mode 100644 index 0000000..afd0a4b --- /dev/null +++ b/adapters/ruby_llm/test/test_evaluation_context.rb @@ -0,0 +1,72 @@ +# 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_finishes_delivery_without_mutating_configuration + subscribe(async: true) + caller_thread = Thread.current + delivery_threads = [] + original = Adapter.reporter.method(:report_now) + Adapter.reporter.define_singleton_method(:report_now) do |trace| + delivery_threads << Thread.current + original.call(trace) + 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_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..25b96bf --- /dev/null +++ b/docs/issues/evaluation-trace-context.md @@ -0,0 +1,16 @@ +# 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 34 tests / 82 assertions and RubyLLM adapter 17 tests / 57 +assertions pass on Ruby 4.0.2. New tests cover actual trace-ID correlation, judge +identity, nested context restoration, synchronous delivery, 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..262c199 --- /dev/null +++ b/docs/milestones/evaluation-trace-context.md @@ -0,0 +1,7 @@ +# 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. +- [ ] Merge the draft PR and release the updated adapter. diff --git a/docs/pull-requests/evaluation-trace-context.md b/docs/pull-requests/evaluation-trace-context.md new file mode 100644 index 0000000..0adeb1e --- /dev/null +++ b/docs/pull-requests/evaluation-trace-context.md @@ -0,0 +1,11 @@ +# 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: 51 tests, 139 assertions, zero failures. Changed +Ruby files pass the Omakase lint configuration. No provider or collector network +requests are made by the tests. Merge/release is intentionally pending review. From ebbfec5c3c81672b016de15cffd40421ef6829fc Mon Sep 17 00:00:00 2001 From: Justin Bowen Date: Sat, 12 Sep 2026 20:55:03 -0700 Subject: [PATCH 2/6] chore(release): 0.3.0 Bumps the core gem and the RubyLLM adapter to 0.3.0 and moves the unreleased changelog entries under that version, so the scoped `with_agent` extension ships as a release consumers can pin instead of a git branch. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01GnvoL2E5F411mxoxTgsn8q --- CHANGELOG.md | 17 ++++++++++++++++- .../activeagents/telemetry/ruby_llm/version.rb | 2 +- lib/activeagents/telemetry/version.rb | 2 +- 3 files changed, 18 insertions(+), 3 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4a5597d..7f81826 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,21 @@ # 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 completed trace before delivery, so the caller can keep the + `trace_id` the collector will store. `synchronous: true` delivers through + `Reporter#report_now` in the calling thread for that scope only; the shared + asynchronous configuration is untouched, and a delivery failure follows the + reporter's existing logging policy rather than raising. 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) ### Fixed 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/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 From 7610bda036b7c1343722af80da14fa544c2a1313 Mon Sep 17 00:00:00 2001 From: Justin Bowen Date: Sat, 12 Sep 2026 21:02:12 -0700 Subject: [PATCH 3/6] fix(ruby_llm): keep a turn's scope and honour sampling on synchronous delivery A turn now snapshots the `with_agent` scope it starts under, so a turn left open by a pending tool call and closed later by `flush!` (or by the next chat) reports with that scope's agent, attributes, callback and delivery mode instead of whatever scope is active at report time. `synchronous: true` no longer routes through `Reporter#report_now`, which skips the sampling and configuration checks. `Reporter#report` takes `sync: true` instead, delivering in the calling thread after the same checks as ordinary delivery, and returns whether the traces were accepted; `BatchingReporter#report` flushes in the calling thread when asked the same. The adapter fires `on_trace` only for an accepted trace, so a caller never keeps the ID of a trace `sample_rate` dropped. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01GnvoL2E5F411mxoxTgsn8q --- CHANGELOG.md | 25 ++++++++---- adapters/ruby_llm/README.md | 19 +++++---- .../activeagents-telemetry-ruby_llm.gemspec | 2 +- .../lib/activeagents/telemetry/ruby_llm.rb | 17 +++++--- .../ruby_llm/test/test_evaluation_context.rb | 40 +++++++++++++++++-- .../telemetry/batching_reporter.rb | 18 +++++---- lib/activeagents/telemetry/reporter.rb | 26 ++++++++---- test/test_parity.rb | 19 ++++++++- test/test_telemetry.rb | 21 +++++++++- 9 files changed, 144 insertions(+), 43 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7f81826..84a8879 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,14 +8,23 @@ `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 completed trace before delivery, so the caller can keep the - `trace_id` the collector will store. `synchronous: true` delivers through - `Reporter#report_now` in the calling thread for that scope only; the shared - asynchronous configuration is untouched, and a delivery failure follows the - reporter's existing logging policy rather than raising. 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) + receives each trace the reporter accepted, so the caller can keep the + `trace_id` the collector will store; a trace dropped by `sample_rate` or a + disabled configuration is never announced. `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. + 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` flushes its buffer in the calling thread when + asked the same. 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 aa01a3f..3d1292f 100644 --- a/adapters/ruby_llm/README.md +++ b/adapters/ruby_llm/README.md @@ -110,14 +110,19 @@ ActiveAgents::Telemetry::RubyLLM.with_agent( end ``` -The callback receives the completed trace before delivery, so an evaluation -result can retain the exact trace ID. `synchronous: true` waits for the existing -reporter's delivery attempt in this scope; it does not mutate the shared async +The callback receives each trace the reporter accepted, so an evaluation result +can retain the exact trace ID that was sent. 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. 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. +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. 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 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 9e4c924..68fcc80 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. @@ -95,6 +95,12 @@ def reporter # 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. `on_trace` runs only for a trace the reporter accepted: one + # dropped by `sample_rate` or a disabled configuration is never + # announced. def with_agent(name, action: "chat", attributes: {}, on_trace: nil, synchronous: false) previous = Thread.current[AGENT_KEY] Thread.current[AGENT_KEY] = { @@ -107,7 +113,7 @@ def with_agent(name, action: "chat", attributes: {}, on_trace: nil, synchronous: 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 @@ -133,6 +139,7 @@ def begin_round(payload) turn = state turn.chat_key = chat_key turn.started_at ||= Time.now + turn.agent ||= Thread.current[AGENT_KEY] end turn.depth += 1 end @@ -167,7 +174,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 || Thread.current[AGENT_KEY] || resolve_agent(payload) || DEFAULT_AGENT started_at = turn.started_at || Time.now finished_at = Time.now error = payload[:exception_object] @@ -212,8 +219,8 @@ def report_turn(payload, turn) trace.add_span(tool_span) end - notify_trace(agent[:on_trace], trace) - agent[:synchronous] ? reporter.report_now(trace) : 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) diff --git a/adapters/ruby_llm/test/test_evaluation_context.rb b/adapters/ruby_llm/test/test_evaluation_context.rb index afd0a4b..800ee3c 100644 --- a/adapters/ruby_llm/test/test_evaluation_context.rb +++ b/adapters/ruby_llm/test/test_evaluation_context.rb @@ -39,14 +39,14 @@ def test_context_is_restored_after_nested_judge_and_exception assert_equal [ "inner", "outer", nil ], roots.map { |root| root["attributes"]["eval.run_id"] } end - def test_synchronous_scope_finishes_delivery_without_mutating_configuration + def test_synchronous_scope_delivers_in_the_calling_thread_without_mutating_configuration subscribe(async: true) caller_thread = Thread.current delivery_threads = [] - original = Adapter.reporter.method(:report_now) - Adapter.reporter.define_singleton_method(:report_now) do |trace| + captured = posted + Adapter.reporter.define_singleton_method(:deliver) do |body| delivery_threads << Thread.current - original.call(trace) + captured << body end Adapter.with_agent("Support", synchronous: true) do @@ -58,6 +58,38 @@ def test_synchronous_scope_finishes_delivery_without_mutating_configuration 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_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 diff --git a/lib/activeagents/telemetry/batching_reporter.rb b/lib/activeagents/telemetry/batching_reporter.rb index 530bd6c..2bc97ae 100644 --- a/lib/activeagents/telemetry/batching_reporter.rb +++ b/lib/activeagents/telemetry/batching_reporter.rb @@ -24,13 +24,16 @@ 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` flushes + # the buffer in the calling thread once the trace is enqueued, so a + # short-lived process can hand a trace over before it exits. + # @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? batch = nil @mutex.synchronize do @@ -38,8 +41,9 @@ def report(traces) batch = @buffer.slice!(0..) if @buffer.size >= configuration.batch_size start_flusher end - deliver_batch(batch) if batch - nil + deliver_batch(batch, blocking: sync) if batch + flush if sync + 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/test/test_parity.rb b/test/test_parity.rb index 1b7cdf8..20664eb 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,23 @@ def test_buffers_until_batch_size_then_delivers_together assert_equal 3, captured.first["traces"].size end + def test_sync_flushes_the_buffer_in_the_calling_thread + reporter = ActiveAgents::Telemetry::BatchingReporter.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 + reporter.shutdown + 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 From 8e6585ea8cb4de545e2ed2af73389a0c974d257e Mon Sep 17 00:00:00 2001 From: Justin Bowen Date: Sat, 12 Sep 2026 21:12:32 -0700 Subject: [PATCH 4/6] fix(ruby_llm): capture a turn's scope once and describe acceptance honestly A turn now captures the `with_agent` scope on its first round whether or not one is active, and report time no longer consults the current scope, so a turn that started unscoped keeps its own identity even when the chat that flushes it runs inside a judge scope. The changelog and README no longer say the callback's trace ID is one the collector will store: the reporter accepted it for delivery, a failed delivery is logged rather than announced, and an asynchronous delivery can outlive a process that exits at once. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01GnvoL2E5F411mxoxTgsn8q --- CHANGELOG.md | 11 +++++---- adapters/ruby_llm/README.md | 23 +++++++++++-------- .../lib/activeagents/telemetry/ruby_llm.rb | 18 ++++++++++----- .../ruby_llm/test/test_evaluation_context.rb | 16 +++++++++++++ 4 files changed, 49 insertions(+), 19 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 84a8879..27c024f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,15 +8,18 @@ `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, so the caller can keep the - `trace_id` the collector will store; a trace dropped by `sample_rate` or a - disabled configuration is never announced. `synchronous: true` delivers in + 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. + 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) diff --git a/adapters/ruby_llm/README.md b/adapters/ruby_llm/README.md index 3d1292f..000dcac 100644 --- a/adapters/ruby_llm/README.md +++ b/adapters/ruby_llm/README.md @@ -110,15 +110,20 @@ ActiveAgents::Telemetry::RubyLLM.with_agent( end ``` -The callback receives each trace the reporter accepted, so an evaluation result -can retain the exact trace ID that was sent. 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. Context +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 diff --git a/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm.rb b/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm.rb index 68fcc80..ebe88b9 100644 --- a/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm.rb +++ b/adapters/ruby_llm/lib/activeagents/telemetry/ruby_llm.rb @@ -98,9 +98,10 @@ def reporter # # 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. `on_trace` runs only for a trace the reporter accepted: one - # dropped by `sample_rate` or a disabled configuration is never - # announced. + # 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] = { @@ -138,8 +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 - turn.agent ||= Thread.current[AGENT_KEY] + 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 @@ -174,7 +180,7 @@ def build_tool_span(payload, started_at, finished_at) private def report_turn(payload, turn) - agent = 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] diff --git a/adapters/ruby_llm/test/test_evaluation_context.rb b/adapters/ruby_llm/test/test_evaluation_context.rb index 800ee3c..7f080e0 100644 --- a/adapters/ruby_llm/test/test_evaluation_context.rb +++ b/adapters/ruby_llm/test/test_evaluation_context.rb @@ -72,6 +72,22 @@ def test_synchronous_scope_still_honours_sampling_and_skips_the_callback 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) From 37404a25bcd65534f1cc4291ffd1967c8d67b9bf Mon Sep 17 00:00:00 2001 From: Justin Bowen Date: Sat, 12 Sep 2026 21:21:33 -0700 Subject: [PATCH 5/6] fix(core): deliver a synchronous batching call's traces without racing the flusher `BatchingReporter#report(sync: true)` appended to the buffer and then flushed after releasing the mutex, so the background flusher could claim the buffer first and the call returned before its own traces went out. A synchronous call now delivers its own traces directly, blocking, and leaves the buffer on its schedule. The branch docs now state the release plan the changelog already records: tag `v0.3.0` after the merge. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01GnvoL2E5F411mxoxTgsn8q --- CHANGELOG.md | 5 +++-- docs/issues/evaluation-trace-context.md | 5 +++-- docs/milestones/evaluation-trace-context.md | 3 ++- docs/pull-requests/evaluation-trace-context.md | 5 +++-- lib/activeagents/telemetry/batching_reporter.rb | 15 ++++++++++----- test/test_parity.rb | 15 +++++++++++---- 6 files changed, 32 insertions(+), 16 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 27c024f..5b51de3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -26,8 +26,9 @@ - `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` flushes its buffer in the calling thread when - asked the same. The RubyLLM adapter now requires core `~> 0.3` for it. + `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/docs/issues/evaluation-trace-context.md b/docs/issues/evaluation-trace-context.md index 25b96bf..1fa833b 100644 --- a/docs/issues/evaluation-trace-context.md +++ b/docs/issues/evaluation-trace-context.md @@ -10,7 +10,8 @@ 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 34 tests / 82 assertions and RubyLLM adapter 17 tests / 57 +Validation: core 38 tests / 93 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, and callback failure. +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 index 262c199..d632483 100644 --- a/docs/milestones/evaluation-trace-context.md +++ b/docs/milestones/evaluation-trace-context.md @@ -4,4 +4,5 @@ - [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. -- [ ] Merge the draft PR and release the updated adapter. +- [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 index 0adeb1e..28610cf 100644 --- a/docs/pull-requests/evaluation-trace-context.md +++ b/docs/pull-requests/evaluation-trace-context.md @@ -6,6 +6,7 @@ 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: 51 tests, 139 assertions, zero failures. Changed +Core and adapter suites pass: 58 tests, 164 assertions, zero failures. Changed Ruby files pass the Omakase lint configuration. No provider or collector network -requests are made by the tests. Merge/release is intentionally pending review. +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 2bc97ae..7600c7b 100644 --- a/lib/activeagents/telemetry/batching_reporter.rb +++ b/lib/activeagents/telemetry/batching_reporter.rb @@ -24,9 +24,10 @@ def initialize(configuration, **options) @shutdown = false end - # Enqueues a trace, flushing if the batch is full. `sync: true` flushes - # the buffer in the calling thread once the trace is enqueued, so a - # short-lived process can hand a trace over before it exits. + # 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 @@ -35,14 +36,18 @@ def report(traces, sync: false) return false if accepted.empty? return false unless configuration.enabled? && configuration.configured? + if sync + deliver_batch(accepted, blocking: true) + return true + end + batch = nil @mutex.synchronize do @buffer.concat(accepted) batch = @buffer.slice!(0..) if @buffer.size >= configuration.batch_size start_flusher end - deliver_batch(batch, blocking: sync) if batch - flush if sync + deliver_batch(batch) if batch true end diff --git a/test/test_parity.rb b/test/test_parity.rb index 20664eb..5a10131 100644 --- a/test/test_parity.rb +++ b/test/test_parity.rb @@ -165,13 +165,20 @@ def test_buffers_until_batch_size_then_delivers_together assert_equal 3, captured.first["traces"].size end - def test_sync_flushes_the_buffer_in_the_calling_thread + 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)) - threads = [] - reporter.define_singleton_method(:deliver) { |_body| threads << Thread.current } + 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 ], threads + 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 From 175a1824dc1da9a52bc1a5b9f0b0ab8df1317a09 Mon Sep 17 00:00:00 2001 From: Justin Bowen Date: Sun, 13 Sep 2026 10:28:41 -0700 Subject: [PATCH 6/6] fix(core): report a synchronous batching call whose payload cannot be built as not accepted `BatchingReporter#report(sync: true)` returned true even when building the payload raised, so the RubyLLM adapter announced a trace that never reached delivery. The synchronous path now builds the payload first and answers false, logging the error, when that fails; a delivery failure after a built payload is still logged and counts as accepted, as in `Reporter#report`. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01E4eTvjQ9tjpDJjHfJNYXXc --- docs/issues/evaluation-trace-context.md | 2 +- docs/pull-requests/evaluation-trace-context.md | 2 +- lib/activeagents/telemetry/batching_reporter.rb | 8 +++++++- test/test_parity.rb | 9 +++++++++ 4 files changed, 18 insertions(+), 3 deletions(-) diff --git a/docs/issues/evaluation-trace-context.md b/docs/issues/evaluation-trace-context.md index 1fa833b..38acc5b 100644 --- a/docs/issues/evaluation-trace-context.md +++ b/docs/issues/evaluation-trace-context.md @@ -10,7 +10,7 @@ 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 38 tests / 93 assertions and RubyLLM adapter 20 tests / 71 +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. diff --git a/docs/pull-requests/evaluation-trace-context.md b/docs/pull-requests/evaluation-trace-context.md index 28610cf..b64dce5 100644 --- a/docs/pull-requests/evaluation-trace-context.md +++ b/docs/pull-requests/evaluation-trace-context.md @@ -6,7 +6,7 @@ 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: 58 tests, 164 assertions, zero failures. Changed +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 7600c7b..757c010 100644 --- a/lib/activeagents/telemetry/batching_reporter.rb +++ b/lib/activeagents/telemetry/batching_reporter.rb @@ -37,7 +37,13 @@ def report(traces, sync: false) return false unless configuration.enabled? && configuration.configured? if sync - deliver_batch(accepted, blocking: true) + 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 diff --git a/test/test_parity.rb b/test/test_parity.rb index 5a10131..c912d3a 100644 --- a/test/test_parity.rb +++ b/test/test_parity.rb @@ -182,6 +182,15 @@ def test_sync_delivers_its_own_traces_in_the_calling_thread_and_leaves_the_buffe 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))