From 9c738a4652bb996debb0aa6d54171174a7ecc395 Mon Sep 17 00:00:00 2001 From: Jan Tychtl Date: Thu, 8 Oct 2026 16:50:43 +0200 Subject: [PATCH 1/5] feat(gooddata-eval): put gen-ai turns under the eval item root trace With GOODDATA_EVAL_JOIN_GENAI_TRACE on, every chat request of an eval item carries W3C baggage naming the item's trace and root span (langfuse_trace_id, langfuse_experiment_item_root_observation_id) and its experiment (langfuse_experiment_id/_name/_dataset_id/_item_id). A gen-ai that accepts caller baggage opens its turn under that root and reports the trace id on response_started. Both ids are derived from the conversation id, so the sending thread and the deferred scorer agree without passing state between threads. For a joined conversation, the deferred task reads the trace by id (/v2/observations?traceId=), waiting until every joined turn is ingested, then exports gd-eval's root into that trace in the sdk-experiment environment, spanning the whole run, and writes each score once on it. Langfuse then reports the item's real cost and latency, and the score rows are halved. Conversations that did not join, and every run with the switch off, keep today's path: poll by session, root in its own trace, scores on both. The item scope is set before the K runs in all twelve agentic kinds. anomaly_detection and what_if now pass conversation_id to observe, so an unlinked root there now carries session.id. jira: trivial risk: low --- .../gooddata_eval/core/agentic/_langfuse.py | 148 +++++++++- .../core/agentic/_trace_linker.py | 65 +++-- .../gooddata_eval/core/agentic/alert_skill.py | 54 ++-- .../core/agentic/anomaly_detection.py | 56 ++-- .../core/agentic/conversation.py | 55 ++-- .../core/agentic/dashboard_skill.py | 54 ++-- .../core/agentic/general_question.py | 52 ++-- .../gooddata_eval/core/agentic/guardrail.py | 52 ++-- .../gooddata_eval/core/agentic/kda_skill.py | 54 ++-- .../core/agentic/metric_skill.py | 54 ++-- .../core/agentic/report_skill.py | 56 ++-- .../gooddata_eval/core/agentic/search_tool.py | 52 ++-- .../core/agentic/visualization.py | 54 ++-- .../src/gooddata_eval/core/agentic/what_if.py | 56 ++-- .../src/gooddata_eval/core/chat/sse_client.py | 43 ++- .../src/gooddata_eval/core/langfuse/client.py | 3 + .../gooddata_eval/core/langfuse/experiment.py | 16 +- .../gooddata_eval/core/langfuse/item_scope.py | 151 ++++++++++ .../core/langfuse/observations.py | 74 +++-- .../src/gooddata_eval/core/models.py | 2 + packages/gooddata-eval/tests/conftest.py | 1 + .../tests/test_agentic_anomaly_detection.py | 2 +- .../tests/test_agentic_join_trace.py | 271 ++++++++++++++++++ .../tests/test_agentic_runner.py | 50 +++- .../tests/test_agentic_what_if.py | 2 +- .../tests/test_langfuse_item_scope.py | 49 ++++ .../tests/test_langfuse_observations.py | 65 ++++- .../gooddata-eval/tests/test_sse_client.py | 113 ++++++++ 28 files changed, 1348 insertions(+), 356 deletions(-) create mode 100644 packages/gooddata-eval/src/gooddata_eval/core/langfuse/item_scope.py create mode 100644 packages/gooddata-eval/tests/test_agentic_join_trace.py create mode 100644 packages/gooddata-eval/tests/test_langfuse_item_scope.py diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py index f6e9ff5b5..e6a011341 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py @@ -10,7 +10,7 @@ from collections.abc import Iterator from contextlib import contextmanager from datetime import datetime, timedelta, timezone -from typing import Any +from typing import TYPE_CHECKING, Any from gooddata_eval.core.agentic._trace_linker import link_cancel_event, linking_is_inline, warn_from_worker from gooddata_eval.core.config import ReasoningEffort, env_flag, normalize_reasoning_effort @@ -22,11 +22,24 @@ ScoreTarget, build_experiment_root_span, ) +from gooddata_eval.core.langfuse.item_scope import ( + SDK_EXPERIMENT_ENVIRONMENT, + ItemTraceScope, + JoinedRun, + joined_run, + joined_turns, + release_joined, + trace_ids_for, +) # Part of this module's public surface: external callers import both names from here. from gooddata_eval.core.langfuse.observations import TraceSummary as _TraceObj # noqa: F401 +from gooddata_eval.core.langfuse.observations import summarize_traces from gooddata_eval.core.langfuse.otlp import Span +if TYPE_CHECKING: + from gooddata_eval.core.agentic._trace_linker import RunIdentity + _log = logging.getLogger(__name__) @@ -200,6 +213,19 @@ def _wait_between_attempts(delay: float) -> bool: return not cancel.is_set() +def _fetch_joined_trace(langfuse: HttpxLangfuseClient, session_id: str, *_window: Any) -> list[Any]: + """The trace a joined conversation's gen-ai turns were put in, read by its id. + + Empty until every joined turn's ``conversation.send_message`` row is ingested: a turn row + lands only when its turn ends, so an early read would cover the first turns alone. + """ + trace_id = trace_ids_for(session_id)[0] + rows = langfuse.list_observations_for_trace(trace_id) + if sum(r.get("name") == "conversation.send_message" for r in rows) < joined_turns(session_id): + return [] + return [t for t in summarize_traces(rows) if t.id == trace_id] + + def find_traces_per_conversation( langfuse: Any, conversation_ids: list[str], @@ -245,12 +271,14 @@ def find_traces_per_conversation( break delay = _INITIAL_DELAY found: list[Any] = [] + joined = joined_run(cid) is not None and isinstance(langfuse, HttpxLangfuseClient) + fetch = _fetch_joined_trace if joined else _fetch_traces_for_session for _attempt in range(_MAX_ATTEMPTS): # Attempt first, sleep only between attempts: a trace that is already ingested # when we look must cost nothing, which is the common case once linking is # batched to the end of the run. try: - found = _fetch_traces_for_session(langfuse, cid, window_start, window_end, pad) + found = fetch(langfuse, cid, window_start, window_end, pad) except Exception as exc: _log.debug("Langfuse trace fetch failed for %s: %s", cid, exc) if found or time.monotonic() + delay > stop_at: @@ -298,8 +326,13 @@ def _experiment_root_span( conversation_id: str | None, item_input: Any, output: Any, + joined_ids: tuple[str, str] | None = None, ) -> Span: - """gd-eval's own root span for one (dataset item, run) -- the whole experiment item.""" + """gd-eval's own root span for one (dataset item, run) -- the whole experiment item. + + ``joined_ids`` is the ``(trace_id, span_id)`` gen-ai's turns were put under; the span then + takes those ids and the environment the langfuse SDK forced on those turns. + """ start, end = _span_window(trace, window) session_id = conversation_id or getattr(trace, "session_id", None) tags = tuple(tag for tag in ("gd-eval", run_metadata.get("testing_framework")) if tag) @@ -328,7 +361,9 @@ def _experiment_root_span( "conversation_id": session_id, }, trace_metadata={"run_name": run_name}, - environment=os.environ.get("LANGFUSE_TRACING_ENVIRONMENT"), + environment=SDK_EXPERIMENT_ENVIRONMENT if joined_ids else os.environ.get("LANGFUSE_TRACING_ENVIRONMENT"), + trace_id=joined_ids[0] if joined_ids else None, + span_id=joined_ids[1] if joined_ids else None, ) @@ -377,6 +412,22 @@ def observe( yield None return + joined = joined_run(conversation_id) + if joined is not None and conversation_id: + yield _export_joined_root( + langfuse, + joined, + conversation_id, + dataset_item_id, + run_name, + run_metadata or {}, + trace=trace, + window=window, + item_input=item_input, + output=output, + ) + return + try: dataset_id = langfuse.dataset_id_for_item(dataset_item_id) except Exception as exc: @@ -431,6 +482,59 @@ def observe( yield ScoreTarget(trace_id, span.trace_id, span.span_id) +def _export_joined_root( + langfuse: HttpxLangfuseClient, + joined: JoinedRun, + conversation_id: str, + dataset_item_id: str, + run_name: str, + run_metadata: dict[str, Any], + *, + trace: Any, + window: tuple[datetime, datetime] | None, + item_input: Any, + output: Any, +) -> ScoreTarget: + """Export the root of a joined conversation's trace and return its one score destination. + + The root takes the trace and span ids gen-ai's turns were sent, and the experiment run + they were told they belong to, so Langfuse assembles one experiment item from both. + ``trace`` is that trace as read before scoring; its window spans every gen-ai row. + """ + trace_id, span_id = trace_ids_for(conversation_id) + if joined.run_name != run_name: + _log.warning( + "Conversation %s joined run %s, not %s; its root follows the join.", + conversation_id, + joined.run_name, + run_name, + ) + span = _experiment_root_span( + trace_id, + dataset_item_id, + joined.dataset_id, + joined.run_name, + run_metadata, + trace=trace, + window=window, + conversation_id=conversation_id, + item_input=item_input, + output=output, + joined_ids=(trace_id, span_id), + ) + try: + langfuse.export_spans([span]) + except Exception as exc: + _log.warning("Failed to export the experiment span for run %s: %s", joined.run_name, exc) + warn_from_worker( + f"[langfuse] WARNING: failed to export experiment span run={joined.run_name} item={dataset_item_id}: {exc}" + ) + return ScoreTarget(trace_id) + finally: + release_joined(conversation_id) + return ScoreTarget(experiment_trace_id=trace_id, experiment_span_id=span_id) + + def score_safe(langfuse: Any, trace_id: Any, **kwargs: Any) -> None: """Create one Langfuse score per destination the target names, ignoring errors. @@ -493,6 +597,42 @@ def log_quality_and_value_scores( ) +def run_context_for(identity: RunIdentity) -> tuple[str, dict[str, Any]]: + """``build_run_context`` for an item's ``RunIdentity``.""" + return build_run_context( + identity.host, + identity.token, + identity.workspace_id, + identity.dataset_name, + identity.run_timestamp, + identity.model_version_override, + identity.run_metadata_extra, + identity.reasoning_effort, + ) + + +def resolve_item_scope( + langfuse: Any, identity: RunIdentity, dataset_item_id: str, *, suffix_runs: bool +) -> ItemTraceScope | None: + """The item's experiment scope, or None when ``observe`` could not export its root. + + That is: a client other than ``HttpxLangfuseClient``, linking switched off, or a dataset + item Langfuse does not know. Sending the join keys then would leave gen-ai's turns under + a root that never arrives. + """ + if not isinstance(langfuse, HttpxLangfuseClient) or env_flag(SKIP_ENV_VAR): + return None + try: + dataset_id = langfuse.dataset_id_for_item(dataset_item_id) + except Exception as exc: + _log.warning("Failed to resolve dataset item %s; its chat turns will not join: %s", dataset_item_id, exc) + return None + if dataset_id is None: + return None + base_name, run_metadata = run_context_for(identity) + return ItemTraceScope(base_name, suffix_runs, dataset_id, dataset_item_id, run_metadata) + + def build_run_context( host: str, token: str, diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_trace_linker.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_trace_linker.py index de77ca445..9907ba4d2 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_trace_linker.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_trace_linker.py @@ -21,6 +21,7 @@ from typing import Any, Protocol from gooddata_eval.core._output import emit_line +from gooddata_eval.core.langfuse.item_scope import ItemTraceScope, join_enabled, release_joined _log = logging.getLogger(__name__) @@ -108,6 +109,23 @@ def open_trace_window(langfuse: Any) -> tuple[Any, datetime]: return (try_make_langfuse_client() if langfuse is None else langfuse), utc_now() +def open_item_trace( + langfuse: Any, identity: RunIdentity, dataset_item_id: str, *, suffix_runs: bool +) -> tuple[Any, datetime, ItemTraceScope | None]: + """``open_trace_window`` plus the item's experiment scope for its chat turns to carry. + + The scope is None unless ``GOODDATA_EVAL_JOIN_GENAI_TRACE`` is on and the item can be + assembled into an experiment; then the run name and dataset are resolved here, before + the run, rather than in the deferred task. + """ + langfuse, window_start = open_trace_window(langfuse) + if not (join_enabled() and langfuse is not None and dataset_item_id): + return langfuse, window_start, None + from gooddata_eval.core.agentic._langfuse import resolve_item_scope # noqa: PLC0415 + + return langfuse, window_start, resolve_item_scope(langfuse, identity, dataset_item_id, suffix_runs=suffix_runs) + + @dataclass(frozen=True) class RunIdentity: """Everything ``build_run_context`` needs to name and describe one eval run. @@ -197,6 +215,7 @@ def submit_trace_scoring( suffix_runs: bool, write_scores: Callable[[RunTraceContext], None], item_input: Any = None, + scope: ItemTraceScope | None = None, ) -> None: """Defer one item's whole Langfuse block: resolve its run context, then write scores. @@ -204,35 +223,35 @@ def submit_trace_scoring( ``find_traces_per_conversation``'s ingestion-lag poll are both round trips publishing an already-decided verdict, so neither belongs on the item's clock. Every caller pins ``window_end`` before calling: a deferred poll must not widen its own query window. + A ``scope`` from ``open_item_trace`` already carries the resolved run context. """ def _link_traces() -> None: from gooddata_eval.core.agentic import _langfuse # noqa: PLC0415 - base_name, run_metadata = _langfuse.build_run_context( - identity.host, - identity.token, - identity.workspace_id, - identity.dataset_name, - identity.run_timestamp, - identity.model_version_override, - identity.run_metadata_extra, - identity.reasoning_effort, - ) - traces = _langfuse.find_traces_per_conversation(langfuse, conversation_ids, window_start, window_end) - write_scores( - RunTraceContext( - run_metadata, - _langfuse, - langfuse, - dataset_item_id, - base_name, - suffix_runs, - traces, - (window_start, window_end), - item_input, + if scope is not None: + base_name, run_metadata = scope.base_name, dict(scope.run_metadata) + else: + base_name, run_metadata = _langfuse.run_context_for(identity) + try: + traces = _langfuse.find_traces_per_conversation(langfuse, conversation_ids, window_start, window_end) + write_scores( + RunTraceContext( + run_metadata, + _langfuse, + langfuse, + dataset_item_id, + base_name, + suffix_runs, + traces, + (window_start, window_end), + item_input, + ) ) - ) + finally: + # Nothing reads a join after this task, including one whose root was never exported. + for conversation_id in conversation_ids: + release_joined(conversation_id) submit_trace_link(_link_traces, item_id=dataset_item_id) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/alert_skill.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/alert_skill.py index 791810386..5e0f83f77 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/alert_skill.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/alert_skill.py @@ -25,7 +25,7 @@ RunIdentity, RunTraceContext, SubmitTraceLink, - open_trace_window, + open_item_trace, run_trace_link_inline, submit_trace_scoring, utc_now, @@ -33,6 +33,7 @@ from gooddata_eval.core.chat.render import render_answer_text from gooddata_eval.core.chat.sse_client import ChatClient, ChatError from gooddata_eval.core.config import ReasoningEffort +from gooddata_eval.core.langfuse.item_scope import item_scope from gooddata_eval.core.models import ( AgenticAssertionError, AgenticEvalOutcome, @@ -904,20 +905,31 @@ def evaluate_agentic_alert_skill( `conversation_id`-on-exception idiom in `ChatClient.ask()`) so callers can retrieve them either way. """ - langfuse, window_start = open_trace_window(langfuse) - summary = run_agentic_alert_skill( - host=host, - token=token, - workspace_id=workspace_id, - question=question, - expected_output=expected_output, - k=k, - max_iterations=max_iterations, - initial_conversation_id=initial_conversation_id, - reasoning_effort=reasoning_effort, - agent_id=agent_id, - user_context=user_context, + identity = RunIdentity( + host, + token, + workspace_id, + dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, ) + langfuse, window_start, scope = open_item_trace(langfuse, identity, dataset_item_id, suffix_runs=k > 1) + with item_scope(scope): + summary = run_agentic_alert_skill( + host=host, + token=token, + workspace_id=workspace_id, + question=question, + expected_output=expected_output, + k=k, + max_iterations=max_iterations, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + user_context=user_context, + ) if langfuse is not None and dataset_item_id: # Pinned on the calling thread: a deferred poll must not widen its query window. @@ -956,23 +968,15 @@ def _write_scores(ctx: RunTraceContext) -> None: # Before the pass@K raise: a failing item's scores are the ones worth having. submit_trace_scoring( submit_trace_link, - RunIdentity( - host, - token, - workspace_id, - dataset_name, - run_timestamp, - model_version_override, - run_metadata_extra, - reasoning_effort, - ), + identity, langfuse=langfuse, dataset_item_id=dataset_item_id, conversation_ids=[r.conversation_id for r in summary.run_results], window_start=window_start, window_end=window_end, - suffix_runs=len(summary.run_results) > 1, + suffix_runs=k > 1, write_scores=_write_scores, + scope=scope, item_input=question, ) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/anomaly_detection.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/anomaly_detection.py index e436efaf6..b980b7067 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/anomaly_detection.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/anomaly_detection.py @@ -44,7 +44,7 @@ RunIdentity, RunTraceContext, SubmitTraceLink, - open_trace_window, + open_item_trace, run_trace_link_inline, submit_trace_scoring, utc_now, @@ -52,6 +52,7 @@ from gooddata_eval.core.chat.render import render_answer_text from gooddata_eval.core.chat.sse_client import ChatClient from gooddata_eval.core.config import ReasoningEffort +from gooddata_eval.core.langfuse.item_scope import item_scope from gooddata_eval.core.models import ( AgenticAssertionError, AgenticEvalOutcome, @@ -530,20 +531,31 @@ def evaluate_agentic_anomaly_detection( user_context: dict | None = None, ) -> AgenticEvalOutcome: """Run anomaly-detection evaluation, log to Langfuse, and raise on failure.""" - langfuse, window_start = open_trace_window(langfuse) - summary = run_agentic_anomaly_detection( - host=host, - token=token, - workspace_id=workspace_id, - question=question, - expected_output=expected_output, - k=k, - max_iterations=max_iterations, - initial_conversation_id=initial_conversation_id, - reasoning_effort=reasoning_effort, - agent_id=agent_id, - user_context=user_context, + identity = RunIdentity( + host, + token, + workspace_id, + dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, ) + langfuse, window_start, scope = open_item_trace(langfuse, identity, dataset_item_id, suffix_runs=k > 1) + with item_scope(scope): + summary = run_agentic_anomaly_detection( + host=host, + token=token, + workspace_id=workspace_id, + question=question, + expected_output=expected_output, + k=k, + max_iterations=max_iterations, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + user_context=user_context, + ) if langfuse is not None and dataset_item_id: # Pinned on the calling thread: a deferred poll must not widen its query window. @@ -574,7 +586,7 @@ def _write_scores(ctx: RunTraceContext) -> None: if name in ev.asserted } ) - with ctx.observe(pt, run_idx) as tid: + with ctx.observe(pt, run_idx, conversation_id=run.conversation_id) as tid: for score_name, value in strict_checks.items(): ctx.score(tid, name=score_name, value=float(value), data_type="BOOLEAN") log_gate_scores(ctx, tid, gate=gate, pass_at_k=summary.pass_at_k, pass_power_k=summary.pass_power_k) @@ -595,23 +607,15 @@ def _write_scores(ctx: RunTraceContext) -> None: # Before the pass@K raise: a failing item's scores are the ones worth having. submit_trace_scoring( submit_trace_link, - RunIdentity( - host, - token, - workspace_id, - dataset_name, - run_timestamp, - model_version_override, - run_metadata_extra, - reasoning_effort, - ), + identity, langfuse=langfuse, dataset_item_id=dataset_item_id, conversation_ids=[r.conversation_id for r in summary.run_results], window_start=window_start, window_end=window_end, - suffix_runs=len(summary.run_results) > 1, + suffix_runs=k > 1, write_scores=_write_scores, + scope=scope, # The question this run answered, so a score is readable without resolving the # conversation back to its item. item_input=question, diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/conversation.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/conversation.py index 25e0260d0..e5f8aef36 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/conversation.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/conversation.py @@ -44,7 +44,7 @@ RunIdentity, RunTraceContext, SubmitTraceLink, - open_trace_window, + open_item_trace, run_trace_link_inline, submit_trace_scoring, utc_now, @@ -56,6 +56,7 @@ from gooddata_eval.core.config import ReasoningEffort from gooddata_eval.core.evaluators._deep_subset import deep_subset from gooddata_eval.core.evaluators._maql import normalize_maql +from gooddata_eval.core.langfuse.item_scope import item_scope from gooddata_eval.core.models import ( AgenticAssertionError, AgenticEvalOutcome, @@ -1094,28 +1095,36 @@ def evaluate_agentic_conversation( either way. """ resolved_mode = resolve_conversation_mode(mode) - langfuse, window_start = open_trace_window(langfuse) - result = run_agentic_conversation( - host=host, - token=token, - workspace_id=workspace_id, - fixture=fixture, - max_clarification_turns=max_clarification_turns, - initial_conversation_id=initial_conversation_id, - reasoning_effort=reasoning_effort, - agent_id=agent_id, - mode=resolved_mode, - clarification_judge=clarification_judge, - user_context=user_context, + identity = RunIdentity( + host, + token, + workspace_id, + dataset_name or fixture.dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, ) + langfuse, window_start, scope = open_item_trace(langfuse, identity, dataset_item_id, suffix_runs=False) + with item_scope(scope): + result = run_agentic_conversation( + host=host, + token=token, + workspace_id=workspace_id, + fixture=fixture, + max_clarification_turns=max_clarification_turns, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + mode=resolved_mode, + clarification_judge=clarification_judge, + user_context=user_context, + ) passed = result.context_success if resolved_mode == "context" else result.conversation_success if langfuse is not None and dataset_item_id: # Pinned on the calling thread: a deferred poll must not widen its query window. window_end = utc_now() - # Resolved here, not inside the task: deferring it would make the queued task hold - # the whole fixture until the drain. - ds_name = dataset_name or fixture.dataset_name failed_turns = {tr.turn_id: tr.failure_reasons() for tr in result.turn_results if not tr.context_success} def _write_scores(ctx: RunTraceContext) -> None: @@ -1203,16 +1212,7 @@ def _write_scores(ctx: RunTraceContext) -> None: # Before the pass@K raise: a failing item's scores are the ones worth having. submit_trace_scoring( submit_trace_link, - RunIdentity( - host, - token, - workspace_id, - ds_name, - run_timestamp, - model_version_override, - run_metadata_extra, - reasoning_effort, - ), + identity, langfuse=langfuse, dataset_item_id=dataset_item_id, conversation_ids=[result.conversation_id], @@ -1220,6 +1220,7 @@ def _write_scores(ctx: RunTraceContext) -> None: window_end=window_end, suffix_runs=False, write_scores=_write_scores, + scope=scope, item_input=fixture.turns[0].message if fixture.turns else fixture.id, ) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/dashboard_skill.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/dashboard_skill.py index c1851383b..ccfe5e2c7 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/dashboard_skill.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/dashboard_skill.py @@ -22,7 +22,7 @@ RunIdentity, RunTraceContext, SubmitTraceLink, - open_trace_window, + open_item_trace, run_trace_link_inline, submit_trace_scoring, utc_now, @@ -30,6 +30,7 @@ from gooddata_eval.core.chat.render import render_answer_text from gooddata_eval.core.chat.sse_client import ChatClient from gooddata_eval.core.config import ReasoningEffort +from gooddata_eval.core.langfuse.item_scope import item_scope from gooddata_eval.core.models import ( AgenticAssertionError, AgenticEvalOutcome, @@ -1057,20 +1058,31 @@ def evaluate_agentic_dashboard_skill( ValueError: the fixture is unusable — see ``_validate_expectation``. Raised before any request, so it means a fixture to fix rather than a result to read. """ - langfuse, window_start = open_trace_window(langfuse) - summary = run_agentic_dashboard_skill( - host=host, - token=token, - workspace_id=workspace_id, - question=question, - expected_output=expected_output, - k=k, - max_iterations=max_iterations, - initial_conversation_id=initial_conversation_id, - reasoning_effort=reasoning_effort, - agent_id=agent_id, - user_context=user_context, + identity = RunIdentity( + host, + token, + workspace_id, + dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, ) + langfuse, window_start, scope = open_item_trace(langfuse, identity, dataset_item_id, suffix_runs=k > 1) + with item_scope(scope): + summary = run_agentic_dashboard_skill( + host=host, + token=token, + workspace_id=workspace_id, + question=question, + expected_output=expected_output, + k=k, + max_iterations=max_iterations, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + user_context=user_context, + ) if langfuse is not None and dataset_item_id: # Pinned on the calling thread: a deferred poll must not widen its query window. @@ -1100,23 +1112,15 @@ def _write_scores(ctx: RunTraceContext) -> None: # Before the pass@K raise: a failing item's scores are the ones worth having. submit_trace_scoring( submit_trace_link, - RunIdentity( - host, - token, - workspace_id, - dataset_name, - run_timestamp, - model_version_override, - run_metadata_extra, - reasoning_effort, - ), + identity, langfuse=langfuse, dataset_item_id=dataset_item_id, conversation_ids=[r.conversation_id for r in summary.run_results], window_start=window_start, window_end=window_end, - suffix_runs=len(summary.run_results) > 1, + suffix_runs=k > 1, write_scores=_write_scores, + scope=scope, item_input=question, ) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/general_question.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/general_question.py index c773bfdcd..8875e6970 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/general_question.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/general_question.py @@ -19,7 +19,7 @@ RunIdentity, RunTraceContext, SubmitTraceLink, - open_trace_window, + open_item_trace, run_trace_link_inline, submit_trace_scoring, utc_now, @@ -28,6 +28,7 @@ from gooddata_eval.core.chat.sse_client import ChatClient from gooddata_eval.core.config import ReasoningEffort from gooddata_eval.core.evaluators._llm_judge import JudgeResponseError, LLMJudge, score_run +from gooddata_eval.core.langfuse.item_scope import item_scope from gooddata_eval.core.models import ( AgenticAssertionError, AgenticEvalOutcome, @@ -279,19 +280,30 @@ def evaluate_agentic_general_question( AgenticEvalOutcome on success; on failure the same three values are attached to the raised exception as ``.reasoning_steps``/``.conversation_id``/``.response_id``. """ - langfuse, window_start = open_trace_window(langfuse) - summary = run_agentic_general_question( - host=host, - token=token, - workspace_id=workspace_id, - question=question, - expected_output=expected_output, - k=k, - initial_conversation_id=initial_conversation_id, - reasoning_effort=reasoning_effort, - agent_id=agent_id, - user_context=user_context, + identity = RunIdentity( + host, + token, + workspace_id, + dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, ) + langfuse, window_start, scope = open_item_trace(langfuse, identity, dataset_item_id, suffix_runs=k > 1) + with item_scope(scope): + summary = run_agentic_general_question( + host=host, + token=token, + workspace_id=workspace_id, + question=question, + expected_output=expected_output, + k=k, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + user_context=user_context, + ) if langfuse is not None and dataset_item_id: # Pinned on the calling thread: a deferred poll must not widen its query window. @@ -322,16 +334,7 @@ def _write_scores(ctx: RunTraceContext) -> None: # Before the pass@K raise: a failing item's scores are the ones worth having. submit_trace_scoring( submit_trace_link, - RunIdentity( - host, - token, - workspace_id, - dataset_name, - run_timestamp, - model_version_override, - run_metadata_extra, - reasoning_effort, - ), + identity, langfuse=langfuse, dataset_item_id=dataset_item_id, # Ungraded runs are skipped above, so polling for their traces would only spend @@ -339,8 +342,9 @@ def _write_scores(ctx: RunTraceContext) -> None: conversation_ids=[r.conversation_id for r in summary.scored_run_results], window_start=window_start, window_end=window_end, - suffix_runs=len(summary.run_results) > 1, + suffix_runs=k > 1, write_scores=_write_scores, + scope=scope, item_input=question, ) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/guardrail.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/guardrail.py index 849a27753..3d771a405 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/guardrail.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/guardrail.py @@ -18,7 +18,7 @@ RunIdentity, RunTraceContext, SubmitTraceLink, - open_trace_window, + open_item_trace, run_trace_link_inline, submit_trace_scoring, utc_now, @@ -28,6 +28,7 @@ from gooddata_eval.core.config import ReasoningEffort from gooddata_eval.core.evaluators._guardrail_criteria import GUARDRAIL_REFUSAL_DEFINITION from gooddata_eval.core.evaluators._llm_judge import JudgeResponseError, LLMJudge, score_run +from gooddata_eval.core.langfuse.item_scope import item_scope from gooddata_eval.core.models import ( AgenticAssertionError, AgenticEvalOutcome, @@ -261,19 +262,30 @@ def evaluate_agentic_guardrail( raised exception as ``.reasoning_steps``/``.conversation_id``/``.response_id`` (mirrors `evaluate_agentic_metric_skill`'s idiom) so callers can retrieve them either way. """ - langfuse, window_start = open_trace_window(langfuse) - summary = run_agentic_guardrail( - host=host, - token=token, - workspace_id=workspace_id, - question=question, - expected_output=expected_output, - k=k, - initial_conversation_id=initial_conversation_id, - reasoning_effort=reasoning_effort, - agent_id=agent_id, - user_context=user_context, + identity = RunIdentity( + host, + token, + workspace_id, + dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, ) + langfuse, window_start, scope = open_item_trace(langfuse, identity, dataset_item_id, suffix_runs=k > 1) + with item_scope(scope): + summary = run_agentic_guardrail( + host=host, + token=token, + workspace_id=workspace_id, + question=question, + expected_output=expected_output, + k=k, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + user_context=user_context, + ) if langfuse is not None and dataset_item_id: # Pinned on the calling thread: a deferred poll must not widen its query window. @@ -304,16 +316,7 @@ def _write_scores(ctx: RunTraceContext) -> None: # Before the pass@K raise: a failing item's scores are the ones worth having. submit_trace_scoring( submit_trace_link, - RunIdentity( - host, - token, - workspace_id, - dataset_name, - run_timestamp, - model_version_override, - run_metadata_extra, - reasoning_effort, - ), + identity, langfuse=langfuse, dataset_item_id=dataset_item_id, # Ungraded runs are skipped above, so polling for their traces would only spend @@ -321,8 +324,9 @@ def _write_scores(ctx: RunTraceContext) -> None: conversation_ids=[r.conversation_id for r in summary.scored_run_results], window_start=window_start, window_end=window_end, - suffix_runs=len(summary.run_results) > 1, + suffix_runs=k > 1, write_scores=_write_scores, + scope=scope, item_input=question, ) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/kda_skill.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/kda_skill.py index fade8bff5..7bf58fea6 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/kda_skill.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/kda_skill.py @@ -22,7 +22,7 @@ RunIdentity, RunTraceContext, SubmitTraceLink, - open_trace_window, + open_item_trace, run_trace_link_inline, submit_trace_scoring, utc_now, @@ -30,6 +30,7 @@ from gooddata_eval.core.chat.render import render_answer_text from gooddata_eval.core.chat.sse_client import ChatClient, ChatError from gooddata_eval.core.config import ReasoningEffort +from gooddata_eval.core.langfuse.item_scope import item_scope from gooddata_eval.core.models import ( AgenticAssertionError, AgenticEvalOutcome, @@ -463,20 +464,31 @@ def evaluate_agentic_kda_skill( AgenticEvalOutcome on success; on failure the same three values are attached to the raised exception as ``.reasoning_steps``/``.conversation_id``/``.response_id``. """ - langfuse, window_start = open_trace_window(langfuse) - summary = run_agentic_kda_skill( - host=host, - token=token, - workspace_id=workspace_id, - question=question, - expected_output=expected_output, - k=k, - max_iterations=max_iterations, - initial_conversation_id=initial_conversation_id, - reasoning_effort=reasoning_effort, - agent_id=agent_id, - user_context=user_context, + identity = RunIdentity( + host, + token, + workspace_id, + dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, ) + langfuse, window_start, scope = open_item_trace(langfuse, identity, dataset_item_id, suffix_runs=k > 1) + with item_scope(scope): + summary = run_agentic_kda_skill( + host=host, + token=token, + workspace_id=workspace_id, + question=question, + expected_output=expected_output, + k=k, + max_iterations=max_iterations, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + user_context=user_context, + ) if langfuse is not None and dataset_item_id: # Pinned on the calling thread: a deferred poll must not widen its query window. @@ -528,23 +540,15 @@ def _write_scores(ctx: RunTraceContext) -> None: # Before the pass@K raise: a failing item's scores are the ones worth having. submit_trace_scoring( submit_trace_link, - RunIdentity( - host, - token, - workspace_id, - dataset_name, - run_timestamp, - model_version_override, - run_metadata_extra, - reasoning_effort, - ), + identity, langfuse=langfuse, dataset_item_id=dataset_item_id, conversation_ids=[r.conversation_id for r in summary.run_results], window_start=window_start, window_end=window_end, - suffix_runs=len(summary.run_results) > 1, + suffix_runs=k > 1, write_scores=_write_scores, + scope=scope, item_input=question, ) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/metric_skill.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/metric_skill.py index ec280c6a9..ebeb8fa2b 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/metric_skill.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/metric_skill.py @@ -25,7 +25,7 @@ RunIdentity, RunTraceContext, SubmitTraceLink, - open_trace_window, + open_item_trace, run_trace_link_inline, submit_trace_scoring, utc_now, @@ -34,6 +34,7 @@ from gooddata_eval.core.chat.sse_client import ChatClient, ChatError from gooddata_eval.core.config import ReasoningEffort from gooddata_eval.core.evaluators._maql import normalize_maql +from gooddata_eval.core.langfuse.item_scope import item_scope from gooddata_eval.core.models import ( AgenticAssertionError, AgenticEvalOutcome, @@ -532,20 +533,31 @@ def evaluate_agentic_metric_skill( `conversation_id`-on-exception idiom in `ChatClient.ask()`) so callers can retrieve them either way. """ - langfuse, window_start = open_trace_window(langfuse) - summary = run_agentic_metric_skill( - host=host, - token=token, - workspace_id=workspace_id, - question=question, - expected_output=expected_output, - k=k, - max_iterations=max_iterations, - initial_conversation_id=initial_conversation_id, - reasoning_effort=reasoning_effort, - agent_id=agent_id, - user_context=user_context, + identity = RunIdentity( + host, + token, + workspace_id, + dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, ) + langfuse, window_start, scope = open_item_trace(langfuse, identity, dataset_item_id, suffix_runs=k > 1) + with item_scope(scope): + summary = run_agentic_metric_skill( + host=host, + token=token, + workspace_id=workspace_id, + question=question, + expected_output=expected_output, + k=k, + max_iterations=max_iterations, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + user_context=user_context, + ) if langfuse is not None and dataset_item_id: # Pinned on the calling thread: a deferred poll must not widen its query window. @@ -577,23 +589,15 @@ def _write_scores(ctx: RunTraceContext) -> None: # Before the pass@K raise: a failing item's scores are the ones worth having. submit_trace_scoring( submit_trace_link, - RunIdentity( - host, - token, - workspace_id, - dataset_name, - run_timestamp, - model_version_override, - run_metadata_extra, - reasoning_effort, - ), + identity, langfuse=langfuse, dataset_item_id=dataset_item_id, conversation_ids=[r.conversation_id for r in summary.run_results], window_start=window_start, window_end=window_end, - suffix_runs=len(summary.run_results) > 1, + suffix_runs=k > 1, write_scores=_write_scores, + scope=scope, item_input=question, ) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/report_skill.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/report_skill.py index 86a62bdb0..0fca3acf4 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/report_skill.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/report_skill.py @@ -23,7 +23,7 @@ RunIdentity, RunTraceContext, SubmitTraceLink, - open_trace_window, + open_item_trace, run_trace_link_inline, submit_trace_scoring, utc_now, @@ -33,6 +33,7 @@ from gooddata_eval.core.chat.sse_client import ChatClient from gooddata_eval.core.config import ReasoningEffort from gooddata_eval.core.evaluators._llm_judge import JudgeResponseError, LLMJudge, score_run +from gooddata_eval.core.langfuse.item_scope import item_scope from gooddata_eval.core.models import ( AgenticAssertionError, AgenticEvalOutcome, @@ -719,21 +720,32 @@ def evaluate_agentic_report_skill( JudgeResponseError: the fixture states a narrative and the judge returned no readable verdict for any run -- an item without a result, not K failures. """ - langfuse, window_start = open_trace_window(langfuse) - summary = run_agentic_report_skill( - host=host, - token=token, - workspace_id=workspace_id, - question=question, - expected_output=expected_output, - k=k, - max_iterations=max_iterations, - initial_conversation_id=initial_conversation_id, - reasoning_effort=reasoning_effort, - agent_id=agent_id, - judge=judge, - user_context=user_context, + identity = RunIdentity( + host, + token, + workspace_id, + dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, ) + langfuse, window_start, scope = open_item_trace(langfuse, identity, dataset_item_id, suffix_runs=k > 1) + with item_scope(scope): + summary = run_agentic_report_skill( + host=host, + token=token, + workspace_id=workspace_id, + question=question, + expected_output=expected_output, + k=k, + max_iterations=max_iterations, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + judge=judge, + user_context=user_context, + ) if langfuse is not None and dataset_item_id: # Pinned on the calling thread: a deferred poll must not widen its query window. @@ -771,23 +783,15 @@ def _write_scores(ctx: RunTraceContext) -> None: # Before the pass@K raise: a failing item's scores are the ones worth having. submit_trace_scoring( submit_trace_link, - RunIdentity( - host, - token, - workspace_id, - dataset_name, - run_timestamp, - model_version_override, - run_metadata_extra, - reasoning_effort, - ), + identity, langfuse=langfuse, dataset_item_id=dataset_item_id, conversation_ids=[r.conversation_id for r in summary.scored_run_results], window_start=window_start, window_end=window_end, - suffix_runs=len(summary.run_results) > 1, + suffix_runs=k > 1, write_scores=_write_scores, + scope=scope, item_input=question, ) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/search_tool.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/search_tool.py index 2adbde990..cdb7256cd 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/search_tool.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/search_tool.py @@ -18,13 +18,14 @@ RunIdentity, RunTraceContext, SubmitTraceLink, - open_trace_window, + open_item_trace, run_trace_link_inline, submit_trace_scoring, utc_now, ) from gooddata_eval.core.chat.sse_client import ChatClient from gooddata_eval.core.config import ReasoningEffort +from gooddata_eval.core.langfuse.item_scope import item_scope from gooddata_eval.core.models import ( AgenticAssertionError, AgenticEvalOutcome, @@ -216,19 +217,30 @@ def evaluate_agentic_search_tool( AgenticEvalOutcome on success; on failure the same three values are attached to the raised exception as ``.reasoning_steps``/``.conversation_id``/``.response_id``. """ - langfuse, window_start = open_trace_window(langfuse) - summary = run_agentic_search_tool( - host=host, - token=token, - workspace_id=workspace_id, - question=question, - expected_tool_call=expected_tool_call, - k=k, - initial_conversation_id=initial_conversation_id, - reasoning_effort=reasoning_effort, - agent_id=agent_id, - user_context=user_context, + identity = RunIdentity( + host, + token, + workspace_id, + dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, ) + langfuse, window_start, scope = open_item_trace(langfuse, identity, dataset_item_id, suffix_runs=k > 1) + with item_scope(scope): + summary = run_agentic_search_tool( + host=host, + token=token, + workspace_id=workspace_id, + question=question, + expected_tool_call=expected_tool_call, + k=k, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + user_context=user_context, + ) if langfuse is not None and dataset_item_id: # Pinned on the calling thread: a deferred poll must not widen its query window. @@ -255,23 +267,15 @@ def _write_scores(ctx: RunTraceContext) -> None: # Before the pass@K raise: a failing item's scores are the ones worth having. submit_trace_scoring( submit_trace_link, - RunIdentity( - host, - token, - workspace_id, - dataset_name, - run_timestamp, - model_version_override, - run_metadata_extra, - reasoning_effort, - ), + identity, langfuse=langfuse, dataset_item_id=dataset_item_id, conversation_ids=[r.conversation_id for r in summary.run_results], window_start=window_start, window_end=window_end, - suffix_runs=len(summary.run_results) > 1, + suffix_runs=k > 1, write_scores=_write_scores, + scope=scope, item_input=question, ) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/visualization.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/visualization.py index 597576001..1285d1c86 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/visualization.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/visualization.py @@ -23,7 +23,7 @@ RunIdentity, RunTraceContext, SubmitTraceLink, - open_trace_window, + open_item_trace, run_trace_link_inline, submit_trace_scoring, utc_now, @@ -38,6 +38,7 @@ evaluation_result_detail, with_execution, ) +from gooddata_eval.core.langfuse.item_scope import item_scope from gooddata_eval.core.models import ( AgenticAssertionError, AgenticEvalOutcome, @@ -430,21 +431,32 @@ def evaluate_agentic_visualization( """ import json as _json # noqa: PLC0415 - langfuse, window_start = open_trace_window(langfuse) - summary = run_agentic_visualization( - host=host, - token=token, - workspace_id=workspace_id, - question=question, - expected_outputs=expected_outputs, - k=k, - max_iterations=max_iterations, - initial_conversation_id=initial_conversation_id, - reasoning_effort=reasoning_effort, - agent_id=agent_id, - user_context=user_context, - requires_execution=requires_execution, + identity = RunIdentity( + host, + token, + workspace_id, + dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, ) + langfuse, window_start, scope = open_item_trace(langfuse, identity, dataset_item_id, suffix_runs=True) + with item_scope(scope): + summary = run_agentic_visualization( + host=host, + token=token, + workspace_id=workspace_id, + question=question, + expected_outputs=expected_outputs, + k=k, + max_iterations=max_iterations, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + user_context=user_context, + requires_execution=requires_execution, + ) if langfuse is not None and dataset_item_id: # Pinned on the calling thread: a deferred poll must not widen its query window. @@ -494,16 +506,7 @@ def _write_scores(ctx: RunTraceContext) -> None: # Before the pass@K raise: a failing item's scores are the ones worth having. submit_trace_scoring( submit_trace_link, - RunIdentity( - host, - token, - workspace_id, - dataset_name, - run_timestamp, - model_version_override, - run_metadata_extra, - reasoning_effort, - ), + identity, langfuse=langfuse, dataset_item_id=dataset_item_id, conversation_ids=[r.conversation_id for r in summary.run_results], @@ -512,6 +515,7 @@ def _write_scores(ctx: RunTraceContext) -> None: # Unlike the other runners, this one suffixes every run, K=1 included. suffix_runs=True, write_scores=_write_scores, + scope=scope, item_input=question, ) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/what_if.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/what_if.py index 3f5f08458..f59954f90 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/what_if.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/what_if.py @@ -38,7 +38,7 @@ RunIdentity, RunTraceContext, SubmitTraceLink, - open_trace_window, + open_item_trace, run_trace_link_inline, submit_trace_scoring, utc_now, @@ -47,6 +47,7 @@ from gooddata_eval.core.chat.sse_client import ChatClient from gooddata_eval.core.config import ReasoningEffort from gooddata_eval.core.evaluators._maql import normalize_maql +from gooddata_eval.core.langfuse.item_scope import item_scope from gooddata_eval.core.models import ( AgenticAssertionError, AgenticEvalOutcome, @@ -491,20 +492,31 @@ def evaluate_agentic_what_if( user_context: dict | None = None, ) -> AgenticEvalOutcome: """Run what-if evaluation, log to Langfuse, and raise WhatIfAssertionError on failure.""" - langfuse, window_start = open_trace_window(langfuse) - summary = run_agentic_what_if( - host=host, - token=token, - workspace_id=workspace_id, - question=question, - expected_output=expected_output, - k=k, - max_iterations=max_iterations, - initial_conversation_id=initial_conversation_id, - reasoning_effort=reasoning_effort, - agent_id=agent_id, - user_context=user_context, + identity = RunIdentity( + host, + token, + workspace_id, + dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, ) + langfuse, window_start, scope = open_item_trace(langfuse, identity, dataset_item_id, suffix_runs=k > 1) + with item_scope(scope): + summary = run_agentic_what_if( + host=host, + token=token, + workspace_id=workspace_id, + question=question, + expected_output=expected_output, + k=k, + max_iterations=max_iterations, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + user_context=user_context, + ) if langfuse is not None and dataset_item_id: # Pinned on the calling thread: a deferred poll must not widen its query window. @@ -537,7 +549,7 @@ def _write_scores(ctx: RunTraceContext) -> None: if name in ev.asserted } ) - with ctx.observe(pt, run_idx) as tid: + with ctx.observe(pt, run_idx, conversation_id=run.conversation_id) as tid: for score_name, value in strict_checks.items(): ctx.score(tid, name=score_name, value=float(value), data_type="BOOLEAN") log_gate_scores(ctx, tid, gate=gate, pass_at_k=summary.pass_at_k, pass_power_k=summary.pass_power_k) @@ -558,23 +570,15 @@ def _write_scores(ctx: RunTraceContext) -> None: # Before the pass@K raise: a failing item's scores are the ones worth having. submit_trace_scoring( submit_trace_link, - RunIdentity( - host, - token, - workspace_id, - dataset_name, - run_timestamp, - model_version_override, - run_metadata_extra, - reasoning_effort, - ), + identity, langfuse=langfuse, dataset_item_id=dataset_item_id, conversation_ids=[r.conversation_id for r in summary.run_results], window_start=window_start, window_end=window_end, - suffix_runs=len(summary.run_results) > 1, + suffix_runs=k > 1, write_scores=_write_scores, + scope=scope, # The question this run answered, so a score is readable without resolving the # conversation back to its item. item_input=question, diff --git a/packages/gooddata-eval/src/gooddata_eval/core/chat/sse_client.py b/packages/gooddata-eval/src/gooddata_eval/core/chat/sse_client.py index e086db8ce..f43661885 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/chat/sse_client.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/chat/sse_client.py @@ -26,6 +26,7 @@ import httpx from gooddata_eval.core.config import ReasoningEffort, normalize_reasoning_effort +from gooddata_eval.core.langfuse.item_scope import JoinedRun, baggage_entries, current_scope, mark_joined, trace_ids_for from gooddata_eval.core.models import ChatResult, DatasetItem _log = logging.getLogger(__name__) @@ -34,6 +35,8 @@ SSE_EVENT_PREFIX = "event: " # gen-ai's last event, only if at least one item was already emitted (conversations_controller.py). _RESPONSE_ENDED_EVENT = "response_ended" +# gen-ai's first event of a turn; its data carries the turn's responseId and traceId. +_RESPONSE_STARTED_EVENT = "response_started" # 500 is here on evidence, not on principle: in one visualization eval batch it hard-failed # 10 of 56 runs with zero retry attempts, and every affected question scored normally when the @@ -247,6 +250,7 @@ class _SseAccumulator: reasoning_steps: list[dict[str, Any]] = field(default_factory=list) adhoc_viz_args: list[dict[str, Any]] = field(default_factory=list) response_id: str | None = None + trace_id: str | None = None stream_ended: bool = False # Reference point for call_ts/result_ts below -- client-observed receipt time, not a # server timestamp, so only meaningful as an offset within this one turn. Wrapped in a @@ -365,6 +369,7 @@ def _build_chat_result(acc: _SseAccumulator) -> ChatResult: } result = ChatResult.model_validate(payload) result.response_id = acc.response_id + result.trace_id = acc.trace_id result.stream_ended = acc.stream_ended return result @@ -456,6 +461,8 @@ def parse_sse_lines(lines: Iterable[str]) -> ChatResult: raise ChatError(message, status_code=code, detail=detail, partial_result=_build_chat_result(acc)) if event_data.get("responseId") and not acc.response_id: acc.response_id = event_data["responseId"] + if current_event == _RESPONSE_STARTED_EVENT and event_data.get("traceId") and not acc.trace_id: + acc.trace_id = event_data["traceId"] item = event_data.get("item") if not item: continue @@ -530,6 +537,9 @@ def __init__( self._reasoning_effort = normalize_reasoning_effort(reasoning_effort) self._agent_id = agent_id self._user_context = user_context + # Run index of each conversation this client has sent to, in first-send order: every + # agentic kind sends run 0 first, then runs 1..K-1 each on a conversation of its own. + self._run_index: dict[str, int] = {} def create_conversation(self) -> str: def _do() -> str: @@ -559,7 +569,8 @@ def send_message( self, conversation_id: str, question: str, *, user_context: dict[str, Any] | None = None ) -> ChatResult: url = f"{self._base}/{conversation_id}/messages" - headers = {**self._auth, "Accept": "text/event-stream", "Content-Type": "application/json"} + join_headers, joined = self._join_item_trace(conversation_id) + headers = {**self._auth, "Accept": "text/event-stream", "Content-Type": "application/json", **join_headers} body: dict[str, Any] = {"item": {"role": "user", "content": {"type": "text", "text": question}}} if self._reasoning_effort is not None: body["options"] = {"reasoningEffort": self._reasoning_effort} @@ -587,7 +598,35 @@ def _do() -> ChatResult: result.turn_wall_clock_sec = time.monotonic() - t0 return result - return _retry_transient(_do, is_retryable=_is_retryable_exc) + try: + result = _retry_transient(_do, is_retryable=_is_retryable_exc) + except ChatError as exc: + if joined is not None and exc.partial_result is not None: + self._record_join(conversation_id, exc.partial_result.trace_id, joined) + raise + if joined is not None: + self._record_join(conversation_id, result.trace_id, joined) + return result + + def _join_item_trace(self, conversation_id: str) -> tuple[dict[str, str], JoinedRun | None]: + """The baggage header carrying the current item's experiment scope, and the run it names. + + Empty and None when no item scope is set, so the request goes out unchanged. + """ + scope = current_scope() + if scope is None: + return {}, None + run_idx = self._run_index.setdefault(conversation_id, len(self._run_index)) + run_name = scope.run_name(run_idx) + labels = [self._auth["baggage"]] if self._auth.get("baggage") else [] + baggage = ",".join([*labels, *baggage_entries(scope, run_name, conversation_id)]) + return {"baggage": baggage}, JoinedRun(run_name, scope.dataset_id) + + @staticmethod + def _record_join(conversation_id: str, reported_trace_id: str | None, run: JoinedRun) -> None: + """Remember the conversation as joined once gen-ai reports the trace it was asked for.""" + if reported_trace_id == trace_ids_for(conversation_id)[0]: + mark_joined(conversation_id, run) def _deadline(self, t0: float) -> tuple[float | None, float, str]: """The earlier of the turn and item caps, as (deadline, budget, scope). diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py index 1441457ec..24a7ca862 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py @@ -198,6 +198,9 @@ def list_traces( self._http, from_time=from_time, to_time=to_time, limit=limit, session_id=session_id ) + def list_observations_for_trace(self, trace_id: str) -> list[dict]: + return observations.list_observations_for_trace(self._http, trace_id) + def flush(self) -> None: pass # no client-side batching diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/experiment.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/experiment.py index b3e9745d6..0383f725d 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/experiment.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/experiment.py @@ -74,15 +74,20 @@ def build_experiment_root_span( observation_metadata: dict[str, Any] | None = None, trace_metadata: dict[str, Any] | None = None, environment: str | None = None, + trace_id: str | None = None, + span_id: str | None = None, ) -> Span: """Build the single root span gd-eval emits per (dataset item, run). With `run=None` this is a plain observation span carrying no `langfuse.experiment.*` attributes at all — used when there is no experiment to attach the item to. + + `trace_id`/`span_id` default to fresh ids; pass them to make the span the root of an + existing trace whose other spans already name it as their parent. """ if end < start: end = start - span_id = new_span_id() + span_id = span_id or new_span_id() attributes: list[dict[str, Any]] = [otlp_attribute(ATTR_OBSERVATION_TYPE, "span")] if item.input is not None: @@ -118,7 +123,14 @@ def build_experiment_root_span( ) attributes.extend(flatten_metadata(ATTR_EXPERIMENT_ITEM_METADATA_PREFIX, item.metadata)) - return Span(trace_id=new_trace_id(), span_id=span_id, name=trace_name, start=start, end=end, attributes=attributes) + return Span( + trace_id=trace_id or new_trace_id(), + span_id=span_id, + name=trace_name, + start=start, + end=end, + attributes=attributes, + ) class ScoreTarget(str): diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/item_scope.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/item_scope.py new file mode 100644 index 000000000..002048b8c --- /dev/null +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/item_scope.py @@ -0,0 +1,151 @@ +# (C) 2026 GoodData Corporation +"""One eval item's Langfuse experiment scope, carried to gen-ai as W3C baggage. + +With ``GOODDATA_EVAL_JOIN_GENAI_TRACE`` on, every chat request of an item tells gen-ai which +trace and parent to open its spans under, and which experiment item they belong to. gen-ai +then builds its turn inside gd-eval's experiment trace instead of a trace of its own, and +gd-eval exports the matching root span afterwards. + +The trace id and the root span id are derived from the conversation id, so the sender and +the deferred scorer (which runs on another thread, with no access to the sender's context) +agree on them without passing anything between threads. +""" + +from __future__ import annotations + +import contextvars +import hashlib +import threading +from collections.abc import Iterator +from contextlib import contextmanager +from dataclasses import dataclass, field +from typing import Any +from urllib.parse import quote + +from gooddata_eval.core.config import env_flag +from gooddata_eval.core.langfuse.experiment import experiment_id_for + +JOIN_ENV_VAR = "GOODDATA_EVAL_JOIN_GENAI_TRACE" + +# The environment the langfuse SDK forces on every span that carries the root-observation +# baggage key. gd-eval's root must use the same one, or Langfuse splits the trace. +SDK_EXPERIMENT_ENVIRONMENT = "sdk-experiment" + +TRACE_ID_KEY = "langfuse_trace_id" +ROOT_OBSERVATION_ID_KEY = "langfuse_experiment_item_root_observation_id" +EXPERIMENT_ID_KEY = "langfuse_experiment_id" +EXPERIMENT_NAME_KEY = "langfuse_experiment_name" +EXPERIMENT_DATASET_ID_KEY = "langfuse_experiment_dataset_id" +EXPERIMENT_ITEM_ID_KEY = "langfuse_experiment_item_id" + + +def join_enabled() -> bool: + return env_flag(JOIN_ENV_VAR) + + +@dataclass(frozen=True) +class ItemTraceScope: + """What every run of one dataset item needs to join its Langfuse experiment. + + ``run_metadata`` rides along so the deferred scorer reuses the name resolved before + the run instead of resolving it again (an unpinned run timestamp would differ). + """ + + base_name: str + suffix_runs: bool + dataset_id: str + item_id: str + run_metadata: dict[str, Any] = field(default_factory=dict, compare=False) + + def run_name(self, run_idx: int) -> str: + """Same rule as ``RunTraceContext.run_name``.""" + return f"{self.base_name}_run{run_idx}" if self.suffix_runs else self.base_name + + +_SCOPE: contextvars.ContextVar[ItemTraceScope | None] = contextvars.ContextVar("gd_eval_item_scope", default=None) + + +@contextmanager +def item_scope(scope: ItemTraceScope | None) -> Iterator[None]: + """Make ``scope`` the current item's scope for the duration of the block. None is a no-op.""" + if scope is None: + yield + return + token = _SCOPE.set(scope) + try: + yield + finally: + _SCOPE.reset(token) + + +def current_scope() -> ItemTraceScope | None: + """The scope of the item running on this thread, or None when joining is off.""" + return _SCOPE.get() if join_enabled() else None + + +def trace_ids_for(conversation_id: str) -> tuple[str, str]: + """``(trace_id, root_span_id)`` for a conversation: 32 and 16 lowercase hex, both nonzero. + + The trace id is the langfuse SDK's ``create_trace_id(seed=conversation_id)``. The span id + is hashed from a different seed, because the SDK's observation-id recipe on the same seed + would just be the trace id's first half. + """ + trace_id = hashlib.sha256(conversation_id.encode("utf-8")).digest()[:16].hex() + span_id = hashlib.sha256(f"{conversation_id}:root".encode()).digest()[:8].hex() + # An all-zero id is invalid in OTel; a sha256 prefix is never zero in practice. + return trace_id if int(trace_id, 16) else "0" * 31 + "1", span_id if int(span_id, 16) else "0" * 15 + "1" + + +def baggage_entries(scope: ItemTraceScope, run_name: str, conversation_id: str) -> list[str]: + """The W3C baggage members that put one run's gen-ai spans under its experiment item.""" + trace_id, span_id = trace_ids_for(conversation_id) + return [ + f"{TRACE_ID_KEY}={trace_id}", + f"{ROOT_OBSERVATION_ID_KEY}={span_id}", + f"{EXPERIMENT_ID_KEY}={experiment_id_for(run_name)}", + f"{EXPERIMENT_NAME_KEY}={quote(run_name, safe='')}", + f"{EXPERIMENT_DATASET_ID_KEY}={quote(scope.dataset_id, safe='')}", + f"{EXPERIMENT_ITEM_ID_KEY}={quote(scope.item_id, safe='')}", + ] + + +@dataclass(frozen=True) +class JoinedRun: + """The experiment a joined conversation's gen-ai spans were told they belong to.""" + + run_name: str + dataset_id: str + + +# Conversations whose gen-ai turns reported the trace id gd-eval asked for, with the number of +# such turns. Process-wide: written by the item thread that sent the turn, read by the +# trace-link thread that scores it, and released once the conversation's root is exported. +_JOINED: dict[str, tuple[JoinedRun, int]] = {} +_JOINED_LOCK = threading.Lock() + + +def mark_joined(conversation_id: str, run: JoinedRun) -> None: + """Record one more turn of ``conversation_id`` that gen-ai put in the requested trace.""" + with _JOINED_LOCK: + turns = _JOINED.get(conversation_id, (run, 0))[1] + _JOINED[conversation_id] = (run, turns + 1) + + +def joined_run(conversation_id: str | None) -> JoinedRun | None: + if not conversation_id: + return None + with _JOINED_LOCK: + entry = _JOINED.get(conversation_id) + return entry[0] if entry else None + + +def joined_turns(conversation_id: str) -> int: + """How many of the conversation's turns gen-ai reported in the requested trace.""" + with _JOINED_LOCK: + entry = _JOINED.get(conversation_id) + return entry[1] if entry else 0 + + +def release_joined(conversation_id: str) -> None: + with _JOINED_LOCK: + _JOINED.pop(conversation_id, None) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/observations.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/observations.py index 51686b8fe..75211ed37 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/observations.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/observations.py @@ -50,29 +50,51 @@ def __init__(self, raw: dict, *, total_cost: float | None = None) -> None: self.end_time: datetime | None = _parse_time(raw.get("endTime")) +# gen-ai's span for one chat turn. In a trace of gen-ai's own it is the parentless root; in +# an eval trace that joined gd-eval's experiment item it sits under gd-eval's root. +_TURN_SPAN_NAME = "conversation.send_message" + + def summarize_traces(rows: list[dict]) -> list[TraceSummary]: """Fold observation rows into one summary per trace, in order of first appearance. - The root is the row without a parent observation; it carries the trace's session, - metadata and latency. Cost is summed over all of the trace's rows, because on a gen-ai - conversation the root has no cost of its own and the model calls under it do. A trace - whose root is not on the page is dropped -- a poll that catches a conversation - mid-ingestion sees children only, and those describe no complete trace. + Session and metadata come from the trace's first gen-ai turn row, else its parentless + row. Latency and the start/end times come from the parentless row; a trace without one + (gen-ai turns under a gd-eval root not exported yet) spans its rows' earliest start to + latest end. Cost is summed over all of the trace's rows, because the model calls carry + it, not the turn or the root. A trace with neither a parentless row nor a turn row is + dropped -- a poll that catches a conversation mid-ingestion sees model calls only, and + those describe no complete trace. """ - order: list[str] = [] - roots: dict[str, dict] = {} - costs: dict[str, float] = {} + grouped: dict[str, list[dict]] = {} for row in rows: trace_id = row.get("traceId") - if not trace_id: - continue - if trace_id not in costs: - order.append(trace_id) - costs[trace_id] = 0.0 - costs[trace_id] += float(row.get("totalCost") or 0.0) - if not row.get("parentObservationId"): - roots.setdefault(trace_id, row) - return [TraceSummary(roots[tid], total_cost=costs[tid]) for tid in order if tid in roots] + if trace_id: + grouped.setdefault(trace_id, []).append(row) + return [summary for trace_rows in grouped.values() if (summary := _fold_trace(trace_rows)) is not None] + + +def _fold_trace(rows: list[dict]) -> TraceSummary | None: + root = next((row for row in rows if not row.get("parentObservationId")), None) + turn = next((row for row in rows if row.get("name") == _TURN_SPAN_NAME), None) + head = turn or root + if head is None: + return None + summary = TraceSummary(head, total_cost=sum(float(row.get("totalCost") or 0.0) for row in rows)) + if root is not None: + summary.root_observation_id = root.get("id") + summary.latency = float(root.get("latency") or 0.0) + summary.start_time = _parse_time(root.get("startTime")) + summary.end_time = _parse_time(root.get("endTime")) + return summary + starts = [t for row in rows if (t := _parse_time(row.get("startTime"))) is not None] + ends = [t for row in rows if (t := _parse_time(row.get("endTime"))) is not None] + summary.root_observation_id = None + summary.start_time = min(starts, default=None) + summary.end_time = max(ends, default=None) + if summary.start_time is not None and summary.end_time is not None: + summary.latency = (summary.end_time - summary.start_time).total_seconds() + return summary def list_traces_in_window( @@ -123,3 +145,21 @@ def list_traces_in_window( if not cursor or len(summaries) >= limit: break return summaries[:limit] + + +def list_observations_for_trace( + http: httpx.Client, trace_id: str, *, page_size: int = 1000, max_pages: int = 8 +) -> list[dict]: + """Every observation row of one trace, following the cursor up to ``max_pages`` pages.""" + params: dict[str, Any] = {"traceId": trace_id, "fields": _FIELDS, "limit": page_size} + rows: list[dict] = [] + cursor: str | None = None + for _page in range(max_pages): + resp = http.get(_OBSERVATIONS_PATH, params=params if cursor is None else {**params, "cursor": cursor}) + resp.raise_for_status() + body = resp.json() + rows.extend(body.get("data") or []) + cursor = (body.get("meta") or {}).get("cursor") + if not cursor: + break + return rows diff --git a/packages/gooddata-eval/src/gooddata_eval/core/models.py b/packages/gooddata-eval/src/gooddata_eval/core/models.py index 1a3c634ef..9b2b1c4a7 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/models.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/models.py @@ -341,6 +341,8 @@ class ChatResult(BaseModel): reasoning_step_events: list[ReasoningStepEvent] = Field(default_factory=list, alias="reasoningStepEvents") conversation_id: str | None = Field(default=None, alias="conversationId") response_id: str | None = Field(default=None, alias="responseId") + # The Langfuse trace gen-ai opened the turn in, from its response_started event. + trace_id: str | None = Field(default=None, alias="traceId") # True once gen-ai's response_ended event arrived. stream_ended: bool = False # Wall-clock seconds for the whole chat turn, timed by the client. diff --git a/packages/gooddata-eval/tests/conftest.py b/packages/gooddata-eval/tests/conftest.py index a0e4e835e..536deb380 100644 --- a/packages/gooddata-eval/tests/conftest.py +++ b/packages/gooddata-eval/tests/conftest.py @@ -25,6 +25,7 @@ def fixtures_dir() -> Path: "GOODDATA_EVAL_CHAT_MAX_BACKOFF_S", "GOODDATA_EVAL_CHAT_TURN_TIMEOUT_S", "GOODDATA_EVAL_CHAT_ITEM_TIMEOUT_S", + "GOODDATA_EVAL_JOIN_GENAI_TRACE", ) diff --git a/packages/gooddata-eval/tests/test_agentic_anomaly_detection.py b/packages/gooddata-eval/tests/test_agentic_anomaly_detection.py index 84e138fb1..84621d158 100644 --- a/packages/gooddata-eval/tests/test_agentic_anomaly_detection.py +++ b/packages/gooddata-eval/tests/test_agentic_anomaly_detection.py @@ -417,7 +417,7 @@ def trace(self, _conversation_id): return None @contextmanager - def observe(self, _trace, _run_idx): + def observe(self, _trace, _run_idx, *, conversation_id=None): yield "trace-id" def score(self, _tid, *, name, value, data_type): diff --git a/packages/gooddata-eval/tests/test_agentic_join_trace.py b/packages/gooddata-eval/tests/test_agentic_join_trace.py new file mode 100644 index 000000000..e0403a10c --- /dev/null +++ b/packages/gooddata-eval/tests/test_agentic_join_trace.py @@ -0,0 +1,271 @@ +# (C) 2026 GoodData Corporation +"""A conversation whose gen-ai turns joined gd-eval's trace is read by trace id and gets its +root exported into that trace, scored once.""" + +from __future__ import annotations + +import json +from contextlib import nullcontext +from datetime import datetime, timezone +from typing import Any +from unittest.mock import patch + +import httpx +import pytest +from gooddata_eval.core.agentic import _langfuse +from gooddata_eval.core.agentic._langfuse import SKIP_ENV_VAR, find_traces_per_conversation, observe, score_safe +from gooddata_eval.core.agentic._trace_linker import RunIdentity, open_item_trace, submit_trace_scoring +from gooddata_eval.core.langfuse.client import HttpxLangfuseClient +from gooddata_eval.core.langfuse.experiment import ScoreTarget, experiment_id_for +from gooddata_eval.core.langfuse.item_scope import ItemTraceScope, JoinedRun, joined_run, mark_joined, trace_ids_for +from gooddata_eval.core.langfuse.otlp import unix_nano + +_OTLP_PATH = "/api/public/otel/v1/traces" +_OBSERVATIONS_PATH = "/api/public/v2/observations" +_WINDOW = (datetime(2026, 9, 8, 9, 0, tzinfo=timezone.utc), datetime(2026, 9, 8, 9, 5, tzinfo=timezone.utc)) +_IDENTITY = RunIdentity("https://h", "tok", "ws", "ds", "2026-09-08", "gpt-x", None, None) + + +@pytest.fixture(autouse=True) +def _langfuse_env(monkeypatch): + monkeypatch.setenv("LANGFUSE_BASE_URL", "https://lf.test") + monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk") + monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk") + monkeypatch.setenv("LANGFUSE_TRACING_ENVIRONMENT", "staging") + monkeypatch.delenv(SKIP_ENV_VAR, raising=False) + + +class _Recorder: + """A Langfuse client over a mock transport; ``rows`` is what the observations read returns.""" + + def __init__(self, rows: list[dict] | None = None) -> None: + self.rows = rows or [] + self.requests: list[httpx.Request] = [] + self.client = HttpxLangfuseClient(transport=httpx.MockTransport(self._handle)) + + def _handle(self, request: httpx.Request) -> httpx.Response: + self.requests.append(request) + if request.url.path.startswith("/api/public/dataset-items/"): + return httpx.Response(200, json={"id": "item-1", "datasetId": "ds-1"}) + if request.url.path == _OBSERVATIONS_PATH: + return httpx.Response(200, json={"data": self.rows, "meta": {}}) + return httpx.Response(200, json={}) + + def spans(self) -> list[dict]: + return [ + span + for request in self.requests + if request.url.path == _OTLP_PATH + for span in json.loads(request.content)["resourceSpans"][0]["scopeSpans"][0]["spans"] + ] + + +def _attrs(span: dict) -> dict[str, Any]: + out: dict[str, Any] = {} + for attr in span["attributes"]: + ((kind, value),) = attr["value"].items() + out[attr["key"]] = [v["stringValue"] for v in value["values"]] if kind == "arrayValue" else value + return out + + +def _turn_rows(conversation_id: str) -> list[dict]: + trace_id, span_id = trace_ids_for(conversation_id) + turn = {"traceId": trace_id, "parentObservationId": span_id, "name": "conversation.send_message"} + return [ + { + **turn, + "id": "t1", + "sessionId": conversation_id, + "startTime": "2026-09-08T09:00:10Z", + "endTime": "2026-09-08T09:00:20Z", + }, + {**turn, "id": "t2", "startTime": "2026-09-08T09:00:30Z", "endTime": "2026-09-08T09:00:50Z"}, + { + "traceId": trace_id, + "id": "g1", + "parentObservationId": "t1", + "totalCost": 0.02, + "startTime": "2026-09-08T09:00:11Z", + "endTime": "2026-09-08T09:00:12Z", + }, + ] + + +def test_a_joined_conversation_is_read_by_its_trace_id_not_polled_by_session(): + cid = "conv-join-read" + mark_joined(cid, JoinedRun("ds_run0", "ds-1")) + rec = _Recorder(_turn_rows(cid)) + + (trace,) = find_traces_per_conversation(rec.client, [cid], *_WINDOW).values() + + (request,) = rec.requests + assert dict(request.url.params)["traceId"] == trace_ids_for(cid)[0] + assert "sessionId" not in request.url.params + assert trace.id == trace_ids_for(cid)[0] + assert trace.latency == 40.0 + assert trace.total_cost == pytest.approx(0.02) + + +def test_a_joined_root_is_exported_into_the_gen_ai_trace_and_scored_once(): + cid = "conv-join-observe" + trace_id, span_id = trace_ids_for(cid) + mark_joined(cid, JoinedRun("ds_run1", "ds-1")) + rec = _Recorder(_turn_rows(cid)) + (trace,) = find_traces_per_conversation(rec.client, [cid], *_WINDOW).values() + + with observe( + rec.client, trace.id, "item-1", "ds_run1", {}, trace=trace, window=_WINDOW, conversation_id=cid, item_input="q" + ) as tid: + score_safe(rec.client, tid, name="passed", value=1.0, data_type="BOOLEAN") + + (span,) = rec.spans() + assert (span["traceId"], span["spanId"]) == (trace_id, span_id) + assert "parentSpanId" not in span + assert span["startTimeUnixNano"] == unix_nano(datetime(2026, 9, 8, 9, 0, 10, tzinfo=timezone.utc)) + assert span["endTimeUnixNano"] == unix_nano(datetime(2026, 9, 8, 9, 0, 50, tzinfo=timezone.utc)) + attrs = _attrs(span) + assert attrs["langfuse.environment"] == "sdk-experiment" + assert attrs["langfuse.experiment.id"] == experiment_id_for("ds_run1") + assert attrs["langfuse.experiment.item.root_observation_id"] == span_id + assert attrs["langfuse.session.id"] == cid + # No dataset lookup: the join already named the dataset. + assert not [r for r in rec.requests if "dataset-items" in r.url.path] + + scores = [json.loads(r.content) for r in rec.requests if r.url.path == "/api/public/scores"] + assert [(s["traceId"], s.get("observationId")) for s in scores] == [(trace_id, span_id)] + assert tid.destinations() == [(trace_id, span_id)] + + +def test_a_joined_root_without_rows_falls_back_to_the_local_window(): + cid = "conv-join-no-rows" + mark_joined(cid, JoinedRun("ds", "ds-1")) + rec = _Recorder() + + with observe(rec.client, None, "item-1", "ds", {}, window=_WINDOW, conversation_id=cid) as tid: + pass + + (span,) = rec.spans() + assert span["traceId"] == trace_ids_for(cid)[0] + assert span["startTimeUnixNano"] == unix_nano(_WINDOW[0]) + assert isinstance(tid, ScoreTarget) + + +def test_an_unjoined_conversation_keeps_its_own_trace_and_both_destinations(): + rec = _Recorder() + with observe( + rec.client, "gen-ai-id", "item-1", "ds", {}, window=_WINDOW, conversation_id="conv-never-joined" + ) as tid: + pass + + (span,) = rec.spans() + assert span["traceId"] != trace_ids_for("conv-never-joined")[0] + assert _attrs(span)["langfuse.environment"] == "staging" + assert len(tid.destinations()) == 2 + + +def test_no_scope_and_no_run_context_lookup_while_the_switch_is_off(monkeypatch): + monkeypatch.delenv("GOODDATA_EVAL_JOIN_GENAI_TRACE", raising=False) + rec = _Recorder() + with patch.object(_langfuse, "build_run_context") as build: + _, _, scope = open_item_trace(rec.client, _IDENTITY, "item-1", suffix_runs=True) + assert scope is None + build.assert_not_called() + assert rec.requests == [] + + +def test_the_scope_is_resolved_before_the_run_when_the_switch_is_on(monkeypatch): + monkeypatch.setenv("GOODDATA_EVAL_JOIN_GENAI_TRACE", "1") + rec = _Recorder() + with patch.object(_langfuse, "build_run_context", return_value=("ds_ts", {"model_version": "m"})): + _, _, scope = open_item_trace(rec.client, _IDENTITY, "item-1", suffix_runs=True) + assert scope == ItemTraceScope("ds_ts", True, "ds-1", "item-1") + assert scope.run_metadata == {"model_version": "m"} + + +@pytest.mark.parametrize("blocker", ["skip", "unknown_item", "foreign_client"]) +def test_no_scope_when_the_root_could_not_be_exported(monkeypatch, blocker): + monkeypatch.setenv("GOODDATA_EVAL_JOIN_GENAI_TRACE", "1") + client: Any = _Recorder().client + if blocker == "skip": + monkeypatch.setenv(SKIP_ENV_VAR, "1") + elif blocker == "unknown_item": + client = HttpxLangfuseClient(transport=httpx.MockTransport(lambda _r: httpx.Response(404))) + else: + client = object() + with patch.object(_langfuse, "build_run_context", return_value=("ds_ts", {})): + assert open_item_trace(client, _IDENTITY, "item-1", suffix_runs=False)[2] is None + + +def test_the_deferred_task_reuses_the_scope_s_run_name(): + scope = ItemTraceScope("ds_pinned", True, "ds-1", "item-1", {"model_version": "m"}) + seen: list[Any] = [] + with ( + patch.object(_langfuse, "build_run_context") as build, + patch.object(_langfuse, "find_traces_per_conversation", return_value={}), + ): + submit_trace_scoring( + lambda task, item_id="": task(), + _IDENTITY, + langfuse=object(), + dataset_item_id="item-1", + conversation_ids=["c"], + window_start=_WINDOW[0], + window_end=_WINDOW[1], + suffix_runs=True, + write_scores=seen.append, + scope=scope, + ) + build.assert_not_called() + (ctx,) = seen + assert ctx.run_name(1) == "ds_pinned_run1" + assert ctx.run_metadata == {"model_version": "m"} + + +def test_a_joined_trace_is_not_read_before_every_sent_turn_is_ingested(monkeypatch): + cid = "conv-join-partial" + for _turn in range(3): + mark_joined(cid, JoinedRun("ds_run0", "ds-1")) + monkeypatch.setattr(_langfuse, "_wait_between_attempts", lambda _delay: False) + rec = _Recorder(_turn_rows(cid)) + + assert find_traces_per_conversation(rec.client, [cid], *_WINDOW) == {cid: None} + + +def test_a_joined_root_releases_its_join_record_once_exported(): + cid = "conv-join-release" + mark_joined(cid, JoinedRun("ds_run0", "ds-1")) + rec = _Recorder(_turn_rows(cid)) + (trace,) = find_traces_per_conversation(rec.client, [cid], *_WINDOW).values() + + with observe(rec.client, trace.id, "item-1", "ds_run0", {}, trace=trace, window=_WINDOW, conversation_id=cid): + pass + + assert joined_run(cid) is None + + +@pytest.mark.parametrize("write_scores_fails", [False, True]) +def test_the_deferred_task_releases_joins_its_scoring_never_exported(write_scores_fails): + cid = f"conv-join-unobserved-{write_scores_fails}" + mark_joined(cid, JoinedRun("ds_run0", "ds-1")) + + def write_scores(_ctx: Any) -> None: + if write_scores_fails: + raise RuntimeError("scoring broke") + + with ( + patch.object(_langfuse, "find_traces_per_conversation", return_value={}), + pytest.raises(RuntimeError) if write_scores_fails else nullcontext(), + ): + submit_trace_scoring( + lambda task, item_id="": task(), + _IDENTITY, + langfuse=object(), + dataset_item_id="item-1", + conversation_ids=[cid], + window_start=_WINDOW[0], + window_end=_WINDOW[1], + suffix_runs=False, + write_scores=write_scores, + ) + + assert joined_run(cid) is None diff --git a/packages/gooddata-eval/tests/test_agentic_runner.py b/packages/gooddata-eval/tests/test_agentic_runner.py index d62ed4f16..2c3db7d08 100644 --- a/packages/gooddata-eval/tests/test_agentic_runner.py +++ b/packages/gooddata-eval/tests/test_agentic_runner.py @@ -20,7 +20,7 @@ ) from gooddata_eval.core.agentic.alert_skill import AlertSkillAssertionError from gooddata_eval.core.evaluators._llm_judge import JudgeResponseError -from gooddata_eval.core.models import AgenticEvalOutcome, DatasetItem +from gooddata_eval.core.models import AgenticEvalOutcome, CreatedVisualization, DatasetItem from gooddata_eval.core.timing import PhaseTimings @@ -968,3 +968,51 @@ def test_run_agentic_items_surfaces_failed_runs_on_a_partial_pass(): report = run_agentic_items([_item()], host="http://host", token="tok", workspace_id="ws1", run_ts="2026-01-01") assert report.items[0].pass_at_k is True assert report.items[0].failed_runs == [_A_FAILED_RUN] + + +class _FirstSend(BaseException): + """Ends a run at its first chat turn; BaseException so no kind's error handling keeps it.""" + + +class _FirstSendClient: + def __init__(self, **_kwargs: Any) -> None: + self.sent: list[str] = [] + + def create_conversation(self) -> str: + return "conv-created" + + def send_message(self, conversation_id: str, *_args: Any, **_kwargs: Any) -> None: + self.sent.append(conversation_id) + raise _FirstSend + + def delete_conversation(self, _conversation_id: str) -> None: + pass + + def close(self) -> None: + pass + + +_RUN_ZERO_KINDS = [ + (m, run_fn, expected or [CreatedVisualization.model_validate(_MIN_VIZ)] if m == "visualization" else expected) + for m, run_fn, _, expected in _CLIENT_CONTEXT_KINDS +] + [("conversation", "run_agentic_conversation", _MIN_CONVERSATION_FIXTURE)] + + +@pytest.mark.parametrize(("module_name", "run_fn", "expected"), _RUN_ZERO_KINDS) +def test_every_kind_sends_run_zero_first( + module_name: str, run_fn: str, expected: Any, monkeypatch: pytest.MonkeyPatch +) -> None: + """ChatClient numbers runs by the order their conversations are first sent to, which + names each run's experiment; that only matches the scorer's run index if run 0 goes first.""" + # The guardrail and general-question kinds build their LLM judge before the first send. + monkeypatch.setenv("OPENAI_API_KEY", "test-key") + module = importlib.import_module(f"gooddata_eval.core.agentic.{module_name}") + client = _FirstSendClient() + if module_name == "conversation": + args: tuple = (module.ConversationFixture.model_validate(expected),) + kwargs: dict[str, Any] = {} + else: + args, kwargs = ("q", expected), {"k": 2} + with patch.object(module, "ChatClient", return_value=client), pytest.raises(_FirstSend): + getattr(module, run_fn)("https://h", "tok", "ws1", *args, initial_conversation_id="conv-run0", **kwargs) + assert client.sent == ["conv-run0"] diff --git a/packages/gooddata-eval/tests/test_agentic_what_if.py b/packages/gooddata-eval/tests/test_agentic_what_if.py index 2f45aea4f..83d93aa70 100644 --- a/packages/gooddata-eval/tests/test_agentic_what_if.py +++ b/packages/gooddata-eval/tests/test_agentic_what_if.py @@ -369,7 +369,7 @@ def trace(self, _conversation_id): return None @contextmanager - def observe(self, _trace, _run_idx): + def observe(self, _trace, _run_idx, *, conversation_id=None): yield "trace-id" def score(self, _tid, *, name, value, data_type): diff --git a/packages/gooddata-eval/tests/test_langfuse_item_scope.py b/packages/gooddata-eval/tests/test_langfuse_item_scope.py new file mode 100644 index 000000000..21e738872 --- /dev/null +++ b/packages/gooddata-eval/tests/test_langfuse_item_scope.py @@ -0,0 +1,49 @@ +# (C) 2026 GoodData Corporation +import hashlib +import re + +import pytest +from gooddata_eval.core.langfuse.item_scope import ( + ItemTraceScope, + current_scope, + item_scope, + trace_ids_for, +) + +_SCOPE = ItemTraceScope(base_name="b", suffix_runs=False, dataset_id="d", item_id="i") + + +def test_trace_id_is_the_langfuse_sdk_seeded_trace_id(): + # langfuse.Langfuse.create_trace_id(seed=...) -- gen-ai and Langfuse tooling agree on it. + assert trace_ids_for("conv-1")[0] == hashlib.sha256(b"conv-1").digest()[:16].hex() + + +def test_ids_are_deterministic_well_formed_and_distinct(): + trace_id, span_id = trace_ids_for("conv-1") + assert trace_ids_for("conv-1") == (trace_id, span_id) + assert re.fullmatch(r"[0-9a-f]{32}", trace_id) + assert re.fullmatch(r"[0-9a-f]{16}", span_id) + assert not trace_id.startswith(span_id) + assert trace_ids_for("conv-2") != (trace_id, span_id) + + +def test_run_name_is_suffixed_only_when_the_item_has_several_runs(): + assert _SCOPE.run_name(0) == "b" + assert ItemTraceScope("b", True, "d", "i").run_name(1) == "b_run1" + + +def test_the_scope_is_visible_only_inside_the_block_and_only_with_the_switch_on(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv("GOODDATA_EVAL_JOIN_GENAI_TRACE", "1") + with item_scope(_SCOPE): + assert current_scope() is _SCOPE + monkeypatch.setenv("GOODDATA_EVAL_JOIN_GENAI_TRACE", "0") + assert current_scope() is None + monkeypatch.setenv("GOODDATA_EVAL_JOIN_GENAI_TRACE", "1") + assert current_scope() is None + + +def test_the_scope_is_reset_when_the_run_raises(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv("GOODDATA_EVAL_JOIN_GENAI_TRACE", "1") + with pytest.raises(RuntimeError), item_scope(_SCOPE): + raise RuntimeError + assert current_scope() is None diff --git a/packages/gooddata-eval/tests/test_langfuse_observations.py b/packages/gooddata-eval/tests/test_langfuse_observations.py index b4c26e6b3..590e067df 100644 --- a/packages/gooddata-eval/tests/test_langfuse_observations.py +++ b/packages/gooddata-eval/tests/test_langfuse_observations.py @@ -5,7 +5,12 @@ import httpx import pytest -from gooddata_eval.core.langfuse.observations import TraceSummary, list_traces_in_window, summarize_traces +from gooddata_eval.core.langfuse.observations import ( + TraceSummary, + list_observations_for_trace, + list_traces_in_window, + summarize_traces, +) def _row( @@ -17,17 +22,21 @@ def _row( total_cost: float | None = None, session_id: str | None = None, metadata: dict | None = None, + name: str | None = None, + start: str = "2026-09-09T10:00:00.000Z", + end: str = "2026-09-09T10:00:12.000Z", ) -> dict: return { "traceId": trace_id, "id": obs_id, + "name": name, "parentObservationId": parent, "latency": latency, "totalCost": total_cost, "sessionId": session_id, "metadata": metadata, - "startTime": "2026-09-09T10:00:00.000Z", - "endTime": "2026-09-09T10:00:12.000Z", + "startTime": start, + "endTime": end, } @@ -70,6 +79,37 @@ def test_a_trace_whose_root_is_not_on_the_page_is_dropped(): assert summarize_traces(rows) == [] +def _turn(obs_id: str, start: str, end: str, **kwargs) -> dict: + return _row("t-1", obs_id, parent="s-root", name="conversation.send_message", start=start, end=end, **kwargs) + + +def test_a_trace_with_gen_ai_turns_but_no_parentless_row_spans_its_turns(): + # A joined trace before gd-eval has exported its root: every gen-ai row has a parent. + rows = [ + _turn("o-t2", "2026-09-09T10:00:20.000Z", "2026-09-09T10:00:30.000Z", session_id="conv-1"), + _turn("o-t1", "2026-09-09T10:00:00.000Z", "2026-09-09T10:00:08.000Z", metadata={"k": "v"}), + _row("t-1", "o-gen", parent="o-t1", total_cost=0.01, end="2026-09-09T10:00:05.000Z"), + ] + (summary,) = summarize_traces(rows) + assert summary.start_time == datetime(2026, 9, 9, 10, 0, 0, tzinfo=timezone.utc) + assert summary.end_time == datetime(2026, 9, 9, 10, 0, 30, tzinfo=timezone.utc) + assert summary.latency == 30.0 + assert summary.total_cost == pytest.approx(0.01) + assert summary.session_id == "conv-1" + assert summary.root_observation_id is None + + +def test_session_and_metadata_come_from_the_gen_ai_turn_not_from_a_gd_eval_root(): + rows = [ + _row("t-1", "s-root", latency=40.0, session_id=None, metadata={"run_name": "r"}, name="gd-eval: q"), + _turn("o-t1", "2026-09-09T10:00:00.000Z", "2026-09-09T10:00:08.000Z", session_id="conv-1", metadata={"k": "v"}), + ] + (summary,) = summarize_traces(rows) + assert (summary.session_id, summary.metadata) == ("conv-1", {"k": "v"}) + assert summary.latency == 40.0 + assert summary.root_observation_id == "s-root" + + def test_the_root_start_and_end_times_are_timezone_aware(): (summary,) = summarize_traces([_row("t-1", "o-root", latency=3.0)]) assert summary.start_time == datetime(2026, 9, 9, 10, 0, 0, tzinfo=timezone.utc) @@ -262,3 +302,22 @@ def handler(request: httpx.Request) -> httpx.Response: now = datetime.now(timezone.utc) with _client(handler) as http, pytest.raises(httpx.HTTPStatusError): list_traces_in_window(http, from_time=now, to_time=now, limit=3, session_id=None) + + +def test_a_trace_read_asks_for_that_trace_and_follows_the_cursor(): + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + if "cursor" not in request.url.params: + return httpx.Response(200, json={"data": [_row("t-1", "o-1")], "meta": {"cursor": "next"}}) + return httpx.Response(200, json={"data": [_row("t-1", "o-2", parent="o-1")], "meta": {}}) + + with _client(handler) as http: + rows = list_observations_for_trace(http, "t-1") + + assert [row["id"] for row in rows] == ["o-1", "o-2"] + assert {r.url.path for r in seen} == {"/api/public/v2/observations"} + params = dict(seen[0].url.params) + assert params == {"traceId": "t-1", "fields": "core,basic,usage,metrics,metadata", "limit": "1000"} + assert seen[1].url.params["cursor"] == "next" diff --git a/packages/gooddata-eval/tests/test_sse_client.py b/packages/gooddata-eval/tests/test_sse_client.py index 3d5d5d574..bfaeffc01 100644 --- a/packages/gooddata-eval/tests/test_sse_client.py +++ b/packages/gooddata-eval/tests/test_sse_client.py @@ -14,6 +14,8 @@ TurnIncompleteError, parse_sse_lines, ) +from gooddata_eval.core.langfuse.experiment import experiment_id_for +from gooddata_eval.core.langfuse.item_scope import ItemTraceScope, JoinedRun, item_scope, joined_run, trace_ids_for from gooddata_eval.core.models import DatasetItem, ReasoningStepEvent, ToolCallEvent, build_latency_breakdown @@ -1130,3 +1132,114 @@ def test_no_trace_labels_send_no_baggage(monkeypatch: pytest.MonkeyPatch, raw: s _client_with_handler(_record_requests(requests)).create_conversation() assert _baggage_of(requests) == [None] + + +def _started_sse(trace_id: str | None) -> bytes: + data = {"responseId": "r1", **({"traceId": trace_id} if trace_id else {})} + return f"event: response_started\ndata: {json.dumps(data)}\n\n".encode() + _OK_SSE + + +def test_parse_sse_lines_reads_the_trace_id_from_response_started(): + lines = ["event: response_started", 'data: {"responseId": "r1", "traceId": "' + "a" * 32 + '"}', ""] + assert parse_sse_lines(lines).trace_id == "a" * 32 + + +def test_parse_sse_lines_ignores_a_trace_id_outside_response_started(): + assert parse_sse_lines(['data: {"traceId": "abc", "responseId": "r1"}']).trace_id is None + + +_SCOPE = ItemTraceScope(base_name="ds_ts_model", suffix_runs=True, dataset_id="ds-1", item_id="item-1") + + +def _scoped_client(monkeypatch: pytest.MonkeyPatch, requests: list[httpx.Request], *, echo_trace: bool = True): + """A client whose fake gen-ai reports, per conversation, the trace id the request asked for.""" + monkeypatch.setenv("GOODDATA_EVAL_JOIN_GENAI_TRACE", "1") + + def handler(request: httpx.Request) -> httpx.Response: + requests.append(request) + asked = (_baggage_of([request])[0] or {}).get("langfuse_trace_id") + return httpx.Response(200, content=_started_sse(asked if echo_trace else "f" * 32)) + + return _client_with_handler(handler) + + +def test_each_conversation_is_sent_its_own_trace_and_root(monkeypatch: pytest.MonkeyPatch) -> None: + requests: list[httpx.Request] = [] + client = _scoped_client(monkeypatch, requests) + with item_scope(_SCOPE): + client.send_message("conv-a-distinct", "q") + client.send_message("conv-b-distinct", "q") + + a, b = _baggage_of(requests) + assert (a["langfuse_trace_id"], a["langfuse_experiment_item_root_observation_id"]) == trace_ids_for( + "conv-a-distinct" + ) + assert a["langfuse_trace_id"] != b["langfuse_trace_id"] + assert a["langfuse_experiment_item_root_observation_id"] != b["langfuse_experiment_item_root_observation_id"] + assert (a["langfuse_experiment_dataset_id"], a["langfuse_experiment_item_id"]) == ("ds-1", "item-1") + + +def test_run_index_follows_the_order_conversations_are_first_sent(monkeypatch: pytest.MonkeyPatch) -> None: + requests: list[httpx.Request] = [] + client = _scoped_client(monkeypatch, requests) + with item_scope(_SCOPE): + # run 0 takes two turns before run 1 starts; its second turn keeps run 0's name. + for conversation_id in ("conv-run0", "conv-run0", "conv-run1"): + client.send_message(conversation_id, "q") + + names = [b["langfuse_experiment_name"] for b in _baggage_of(requests)] + assert names == ["ds_ts_model_run0", "ds_ts_model_run0", "ds_ts_model_run1"] + ids = [b["langfuse_experiment_id"] for b in _baggage_of(requests)] + assert ids == [experiment_id_for(name) for name in names] + + +def test_labels_and_scope_share_one_baggage_header(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("GOODDATA_EVAL_TRACE_LABELS", "github_run_id=42") + requests: list[httpx.Request] = [] + client = _scoped_client(monkeypatch, requests) + with item_scope(_SCOPE): + client.send_message("conv-labels", "q") + + (baggage,) = _baggage_of(requests) + assert baggage["langfuse_metadata_github_run_id"] == "42" + assert "langfuse_trace_id" in baggage + + +def test_a_conversation_is_joined_only_when_gen_ai_reports_the_trace_it_was_asked_for( + monkeypatch: pytest.MonkeyPatch, +) -> None: + client = _scoped_client(monkeypatch, []) + with item_scope(_SCOPE): + result = client.send_message("conv-joined", "q") + assert result.trace_id == trace_ids_for("conv-joined")[0] + assert joined_run("conv-joined") == JoinedRun("ds_ts_model_run0", "ds-1") + + other = _scoped_client(monkeypatch, [], echo_trace=False) + with item_scope(_SCOPE): + other.send_message("conv-not-joined", "q") + assert joined_run("conv-not-joined") is None + + +@pytest.mark.parametrize("switch", [None, "", "0", "false"]) +def test_switch_off_sends_no_new_baggage_even_inside_a_scope(monkeypatch: pytest.MonkeyPatch, switch) -> None: + if switch is None: + monkeypatch.delenv("GOODDATA_EVAL_JOIN_GENAI_TRACE", raising=False) + else: + monkeypatch.setenv("GOODDATA_EVAL_JOIN_GENAI_TRACE", switch) + monkeypatch.delenv("GOODDATA_EVAL_TRACE_LABELS", raising=False) + requests: list[httpx.Request] = [] + client = _client_with_handler(_record_requests(requests)) + with item_scope(_SCOPE): + client.send_message("conv-off", "q") + + assert _baggage_of(requests) == [None] + assert joined_run("conv-off") is None + + +def test_switch_on_without_a_scope_sends_no_new_baggage(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("GOODDATA_EVAL_JOIN_GENAI_TRACE", "1") + monkeypatch.delenv("GOODDATA_EVAL_TRACE_LABELS", raising=False) + requests: list[httpx.Request] = [] + _client_with_handler(_record_requests(requests)).send_message("conv-unscoped", "q") + + assert _baggage_of(requests) == [None] From 9e4d9a42da52c1ed74e6f5decbd16a6f789d29cf Mon Sep 17 00:00:00 2001 From: Jan Tychtl Date: Thu, 8 Oct 2026 16:51:07 +0200 Subject: [PATCH 2/5] feat(gooddata-eval): let a ChatClient subclass read a turn's SSE lines ChatClient.send_message passes each attempt's lines through _tap (identity by default) before parsing. A subclass that needs more of the stream than ChatResult keeps, such as tavern's obfuscation client reading an error's reason, overrides _tap instead of copying send_message, and so keeps the turn deadline and the eval-trace join. jira: trivial risk: low --- .../gooddata-eval/src/gooddata_eval/core/chat/sse_client.py | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/chat/sse_client.py b/packages/gooddata-eval/src/gooddata_eval/core/chat/sse_client.py index f43661885..32ff0fc7c 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/chat/sse_client.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/chat/sse_client.py @@ -590,7 +590,7 @@ def _do() -> ChatResult: with self._client.stream("POST", url, json=body, headers=headers) as resp: resp.raise_for_status() try: - result = parse_sse_lines(_until_deadline(resp.iter_lines(), deadline, budget, scope)) + result = parse_sse_lines(self._tap(_until_deadline(resp.iter_lines(), deadline, budget, scope))) except ChatError as exc: if exc.partial_result is not None: exc.partial_result.turn_wall_clock_sec = time.monotonic() - t0 @@ -608,6 +608,10 @@ def _do() -> ChatResult: self._record_join(conversation_id, result.trace_id, joined) return result + def _tap(self, lines: Iterable[str]) -> Iterable[str]: + """Hook over one attempt's SSE lines, before they are parsed. Identity here.""" + return lines + def _join_item_trace(self, conversation_id: str) -> tuple[dict[str, str], JoinedRun | None]: """The baggage header carrying the current item's experiment scope, and the run it names. From 9f7d4c94e9fd70cf8a3e082100c4df8e1cd969c0 Mon Sep 17 00:00:00 2001 From: Jan Tychtl Date: Thu, 8 Oct 2026 17:01:17 +0200 Subject: [PATCH 3/5] fix(gooddata-eval): wait out the Langfuse rate-limit window on a 429 Langfuse Cloud rate-limits public API calls per organization in fixed one-minute windows, and score writes share the bucket with observation and score reads (30/min on Hobby, 1,000/min on Pro). A 429 carries a Retry-After counting down to the window's reset, often 30-60 s. create_score capped every retry wait at 5 s, so both retries landed in the same exhausted window and the score was dropped. On a Hobby project a 60-request burst lost 43 scores this way. A 429 now waits up to the 60 s window; 5xx retries keep the 5 s cap. jira: trivial risk: low --- .../src/gooddata_eval/core/langfuse/client.py | 9 ++++++--- .../tests/test_langfuse_client.py | 18 +++++++++++++----- 2 files changed, 19 insertions(+), 8 deletions(-) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py index 24a7ca862..1d4817bde 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py @@ -29,6 +29,9 @@ _MAX_SCORE_ATTEMPTS = 3 _DEFAULT_RETRY_DELAY = 0.5 _MAX_RETRY_DELAY = 5.0 +# Langfuse Cloud rate-limits in fixed one-minute windows, and a 429's Retry-After counts down to +# the window's reset. A shorter wait retries into the same exhausted window and loses the score. +_MAX_THROTTLE_DELAY = 60.0 def _is_retryable(resp: httpx.Response) -> bool: @@ -38,8 +41,8 @@ def _is_retryable(resp: httpx.Response) -> bool: def _retry_delay(resp: httpx.Response) -> float: """Seconds to wait before the next attempt, from `Retry-After` when the server names one. - Unparsable, negative or NaN values fall back to the default; the cap bounds the wait so a - throttled score cannot hold a linking worker for long. + Unparsable, negative or NaN values fall back to the default. A 429 waits up to one rate-limit + window; any other retryable status is capped at a few seconds. """ try: asked_for = float(resp.headers.get("Retry-After", "")) @@ -47,7 +50,7 @@ def _retry_delay(resp: httpx.Response) -> float: return _DEFAULT_RETRY_DELAY if not asked_for >= 0: return _DEFAULT_RETRY_DELAY - return min(asked_for, _MAX_RETRY_DELAY) + return min(asked_for, _MAX_THROTTLE_DELAY if resp.status_code == 429 else _MAX_RETRY_DELAY) class _TraceListResult: diff --git a/packages/gooddata-eval/tests/test_langfuse_client.py b/packages/gooddata-eval/tests/test_langfuse_client.py index 931d7588e..74ae5b8d0 100644 --- a/packages/gooddata-eval/tests/test_langfuse_client.py +++ b/packages/gooddata-eval/tests/test_langfuse_client.py @@ -151,14 +151,22 @@ def test_a_retry_after_that_cannot_be_slept_falls_back_to_the_default(header): assert client_module._retry_delay(httpx.Response(429, headers={"Retry-After": header})) == 0.5 -def test_a_retry_after_beyond_the_cap_is_clamped(): - assert client_module._retry_delay(httpx.Response(429, headers={"Retry-After": "60"})) == 5.0 +def test_a_throttled_score_waits_out_the_whole_rate_limit_window(): + # Langfuse limits in fixed one-minute windows and names the wait until the window resets. + assert client_module._retry_delay(httpx.Response(429, headers={"Retry-After": "56"})) == 56.0 + + +def test_a_retry_after_beyond_one_window_is_clamped_to_one_window(): + assert client_module._retry_delay(httpx.Response(429, headers={"Retry-After": "600"})) == 60.0 + + +def test_a_server_error_retry_after_stays_short(): + assert client_module._retry_delay(httpx.Response(503, headers={"Retry-After": "60"})) == 5.0 def test_a_retry_after_given_as_a_date_falls_back_to_the_default(): - # Retry-After is allowed to be an HTTP-date; this client reads seconds only. Langfuse - # documents the header as a number of seconds, and the cap already bounds the wait, so - # parsing a date could only turn the 0.5s fallback into the same 5s ceiling. + # Retry-After is allowed to be an HTTP-date; this client reads seconds only, which is the + # form Langfuse documents and sends. response = httpx.Response(429, headers={"Retry-After": "Wed, 21 Oct 2026 07:28:00 GMT"}) assert client_module._retry_delay(response) == 0.5 From e7c3a7337818eea8ac2331f2bb2ce58d23877e93 Mon Sep 17 00:00:00 2001 From: Jan Tychtl Date: Fri, 9 Oct 2026 12:18:39 +0200 Subject: [PATCH 4/5] feat(gooddata-eval): send an item's scores as one Langfuse request Each eval item wrote every score with its own POST /api/public/scores, twice when the item was not joined to gen-ai's trace. Those requests share Langfuse's public-API rate-limit bucket with the observation and score reads, so a nightly run spent most of its budget on score writes. submit_trace_scoring now collects every score_safe call of the deferred scoring step and sends them as POST /api/public/scores arrays of up to 100. Score ids are fixed before the first attempt, so a retried batch cannot write a score twice. Langfuse validates array entries one by one and answers 207 when some are rejected; that is logged, not retried. Callers outside submit_trace_scoring (the single-shot sink, any client without create_scores) still write each score as it comes. Under a sustained 429 one item now waits at most 2 x 60 s, where per-score writes with the 60 s Retry-After cap could wait 36 minutes, past tavern's 420 s test timeout. A cancelled drain sends what was collected before the interrupt once, without waiting out a retry. jira: trivial risk: low --- .../gooddata_eval/core/agentic/_langfuse.py | 35 ++++++ .../core/agentic/_trace_linker.py | 25 ++--- .../src/gooddata_eval/core/langfuse/client.py | 91 +++++++++++++--- .../tests/test_agentic_join_trace.py | 97 ++++++++++++++++- .../tests/test_langfuse_client.py | 101 ++++++++++++++++++ .../tests/test_langfuse_e2e_fake_server.py | 18 ++-- 6 files changed, 331 insertions(+), 36 deletions(-) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py index e6a011341..e4c9e0a58 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py @@ -3,6 +3,7 @@ from __future__ import annotations +import contextvars import logging import os import threading @@ -535,6 +536,36 @@ def _export_joined_root( return ScoreTarget(experiment_trace_id=trace_id, experiment_span_id=span_id) +# Scores held by an open ``collect_scores`` block on this thread, sent together when it ends. +_PENDING_SCORES: contextvars.ContextVar[list[dict[str, Any]] | None] = contextvars.ContextVar( + "gd_eval_pending_scores", default=None +) + + +@contextmanager +def collect_scores(langfuse: Any) -> Iterator[None]: + """Hold every ``score_safe`` write made in the block and send them as one batch at its end. + + The batch is sent even when the drain was cancelled meanwhile, without retrying, so the + scores made before the interrupt are kept. A client without ``create_scores`` writes each + score as it comes. + """ + if not hasattr(langfuse, "create_scores"): + yield + return + pending: list[dict[str, Any]] = [] + token = _PENDING_SCORES.set(pending) + try: + yield + finally: + _PENDING_SCORES.reset(token) + if pending: + try: + langfuse.create_scores(pending, wait=_wait_between_attempts) + except Exception as exc: + _log.warning("Failed to log %d scores: %s", len(pending), exc) + + def score_safe(langfuse: Any, trace_id: Any, **kwargs: Any) -> None: """Create one Langfuse score per destination the target names, ignoring errors. @@ -548,7 +579,11 @@ def score_safe(langfuse: Any, trace_id: Any, **kwargs: Any) -> None: if _drain_is_cancelled(): return targets = trace_id.destinations() if isinstance(trace_id, ScoreTarget) else [(str(trace_id), None)] + pending = _PENDING_SCORES.get() for target_id, observation_id in targets: + if pending is not None: + pending.append({"trace_id": target_id, "observation_id": observation_id, **kwargs}) + continue extra = {"observation_id": observation_id} if observation_id else {} try: langfuse.create_score(trace_id=target_id, **kwargs, **extra) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_trace_linker.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_trace_linker.py index 9907ba4d2..38bd56913 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_trace_linker.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_trace_linker.py @@ -235,19 +235,20 @@ def _link_traces() -> None: base_name, run_metadata = _langfuse.run_context_for(identity) try: traces = _langfuse.find_traces_per_conversation(langfuse, conversation_ids, window_start, window_end) - write_scores( - RunTraceContext( - run_metadata, - _langfuse, - langfuse, - dataset_item_id, - base_name, - suffix_runs, - traces, - (window_start, window_end), - item_input, + with _langfuse.collect_scores(langfuse): + write_scores( + RunTraceContext( + run_metadata, + _langfuse, + langfuse, + dataset_item_id, + base_name, + suffix_runs, + traces, + (window_start, window_end), + item_input, + ) ) - ) finally: # Nothing reads a join after this task, including one whose root was never exported. for conversation_id in conversation_ids: diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py index 1d4817bde..f16db9209 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py @@ -6,9 +6,12 @@ from __future__ import annotations +import json +import logging import threading import time import uuid +from collections.abc import Callable from datetime import datetime, timezone from typing import Any @@ -26,7 +29,11 @@ _SCORES_PATH = "/api/public/scores" _OTLP_PATH = "/api/public/otel/v1/traces" +_log = logging.getLogger(__name__) + _MAX_SCORE_ATTEMPTS = 3 +# Scores per POST /api/public/scores array; the rate limit counts a request, not its scores. +_SCORE_BATCH_SIZE = 100 _DEFAULT_RETRY_DELAY = 0.5 _MAX_RETRY_DELAY = 5.0 # Langfuse Cloud rate-limits in fixed one-minute windows, and a 429's Retry-After counts down to @@ -53,6 +60,35 @@ def _retry_delay(resp: httpx.Response) -> float: return min(asked_for, _MAX_THROTTLE_DELAY if resp.status_code == 429 else _MAX_RETRY_DELAY) +def _sleep(delay: float) -> bool: + time.sleep(delay) + return True + + +def _score_body( + trace_id: str, + name: str, + value: float, + data_type: str, + comment: str | None = None, + observation_id: str | None = None, +) -> dict[str, Any]: + body: dict[str, Any] = { + "id": str(uuid.uuid4()), + "traceId": trace_id, + "name": name, + # A BOOLEAN score goes over the wire as 1.0/0.0 whatever its Python type: the + # sink's compute_scores yields int 1/0 and the agentic path float 1.0/0.0. + "value": (1.0 if value else 0.0) if data_type == "BOOLEAN" else value, + "dataType": data_type, + } + if comment: + body["comment"] = comment + if observation_id: + body["observationId"] = observation_id + return body + + class _TraceListResult: def __init__(self, data: list[TraceSummary]) -> None: self.data = data @@ -145,26 +181,51 @@ def create_score( observation_id: str | None = None, ) -> None: """Attach one score to a trace, or to a single observation inside it.""" - body: dict[str, Any] = { - "id": str(uuid.uuid4()), - "traceId": trace_id, - "name": name, - # A BOOLEAN score goes over the wire as 1.0/0.0 whatever its Python type: the - # sink's compute_scores yields int 1/0 and the agentic path float 1.0/0.0. - "value": (1.0 if value else 0.0) if data_type == "BOOLEAN" else value, - "dataType": data_type, - } - if comment: - body["comment"] = comment - if observation_id: - body["observationId"] = observation_id + self._post_scores(_score_body(trace_id, name, value, data_type, comment, observation_id)) + + def create_scores(self, scores: list[dict[str, Any]], *, wait: Callable[[float], bool] = _sleep) -> None: + """Send many scores as ``POST /api/public/scores`` arrays of up to ``_SCORE_BATCH_SIZE``. + + Each entry takes ``create_score``'s keyword arguments. An entry that cannot be sent (an + unknown argument, a non-finite value) is logged and skipped, so it does not take the + rest of its batch down with it. A batch Langfuse accepts only in part (207) is logged + and not retried: resending it would duplicate the accepted scores. ``wait`` serves each + retry delay; returning False gives up on the retry. + """ + bodies: list[dict[str, Any]] = [] + for score in scores: + try: + body = _score_body(**score) + # httpx serialises with allow_nan=False; checked here so one entry fails alone. + json.dumps(body, allow_nan=False) + except (TypeError, ValueError) as exc: + _log.warning("Langfuse: skipping score %s: %s", score.get("name"), exc) + continue + bodies.append(body) + for start in range(0, len(bodies), _SCORE_BATCH_SIZE): + batch = bodies[start : start + _SCORE_BATCH_SIZE] + resp = self._post_scores(batch, wait) + if resp.status_code == 207: + result = resp.json() + _log.warning( + "Langfuse: %s of %d scores rejected: %s", + result.get("rejected"), + len(batch), + result.get("errors"), + ) + + def _post_scores( + self, body: dict[str, Any] | list[dict[str, Any]], wait: Callable[[float], bool] = _sleep + ) -> httpx.Response: + # The score ids are fixed before the first attempt, so a retried request cannot + # write a score twice. resp = self._http.post(_SCORES_PATH, json=body) for _retry in range(_MAX_SCORE_ATTEMPTS - 1): - if not _is_retryable(resp): + if not _is_retryable(resp) or not wait(_retry_delay(resp)): break - time.sleep(_retry_delay(resp)) resp = self._http.post(_SCORES_PATH, json=body) resp.raise_for_status() + return resp def export_spans(self, spans: list[otlp.Span]) -> None: """Export spans to Langfuse over OTLP/HTTP JSON. Raises on a refused or rejected export.""" diff --git a/packages/gooddata-eval/tests/test_agentic_join_trace.py b/packages/gooddata-eval/tests/test_agentic_join_trace.py index e0403a10c..969de2d5c 100644 --- a/packages/gooddata-eval/tests/test_agentic_join_trace.py +++ b/packages/gooddata-eval/tests/test_agentic_join_trace.py @@ -5,6 +5,7 @@ from __future__ import annotations import json +import threading from contextlib import nullcontext from datetime import datetime, timezone from typing import Any @@ -13,8 +14,14 @@ import httpx import pytest from gooddata_eval.core.agentic import _langfuse -from gooddata_eval.core.agentic._langfuse import SKIP_ENV_VAR, find_traces_per_conversation, observe, score_safe -from gooddata_eval.core.agentic._trace_linker import RunIdentity, open_item_trace, submit_trace_scoring +from gooddata_eval.core.agentic._langfuse import ( + SKIP_ENV_VAR, + collect_scores, + find_traces_per_conversation, + observe, + score_safe, +) +from gooddata_eval.core.agentic._trace_linker import _CANCEL, RunIdentity, open_item_trace, submit_trace_scoring from gooddata_eval.core.langfuse.client import HttpxLangfuseClient from gooddata_eval.core.langfuse.experiment import ScoreTarget, experiment_id_for from gooddata_eval.core.langfuse.item_scope import ItemTraceScope, JoinedRun, joined_run, mark_joined, trace_ids_for @@ -269,3 +276,89 @@ def write_scores(_ctx: Any) -> None: ) assert joined_run(cid) is None + + +def _score_bodies(rec: _Recorder) -> list[list[dict]]: + return [json.loads(r.content) for r in rec.requests if r.url.path == "/api/public/scores"] + + +def test_scores_collected_for_an_item_go_out_in_one_request(): + rec = _Recorder() + target = ScoreTarget("genai-trace", "eval-trace", "eval-root") + + with collect_scores(rec.client): + score_safe(rec.client, target, name="gate_passed", value=1.0, data_type="BOOLEAN") + score_safe(rec.client, target, name="quality_score", value=0.5, data_type="NUMERIC") + assert _score_bodies(rec) == [] + + (bodies,) = _score_bodies(rec) + assert sorted((b["traceId"], b["name"], b.get("observationId")) for b in bodies) == [ + ("eval-trace", "gate_passed", "eval-root"), + ("eval-trace", "quality_score", "eval-root"), + ("genai-trace", "gate_passed", None), + ("genai-trace", "quality_score", None), + ] + + +def test_scores_collected_before_a_drain_is_cancelled_are_still_sent_once(): + rec = _Recorder() + cancel = threading.Event() + token = _CANCEL.set(cancel) + try: + with collect_scores(rec.client): + score_safe(rec.client, "t-1", name="gate_passed", value=1.0, data_type="BOOLEAN") + cancel.set() + score_safe(rec.client, "t-1", name="quality_score", value=0.5, data_type="NUMERIC") + finally: + _CANCEL.reset(token) + + assert [[b["name"] for b in bodies] for bodies in _score_bodies(rec)] == [["gate_passed"]] + + +def test_a_throttled_flush_in_a_cancelled_drain_gives_up_without_waiting(monkeypatch): + slept: list[float] = [] + monkeypatch.setattr(_langfuse.time, "sleep", slept.append) + calls: list[httpx.Request] = [] + + def throttled(request: httpx.Request) -> httpx.Response: + calls.append(request) + return httpx.Response(429, headers={"Retry-After": "60"}, json={}) + + client = HttpxLangfuseClient(transport=httpx.MockTransport(throttled)) + cancel = threading.Event() + token = _CANCEL.set(cancel) + try: + with collect_scores(client): + score_safe(client, "t-1", name="gate_passed", value=1.0, data_type="BOOLEAN") + cancel.set() + finally: + _CANCEL.reset(token) + + assert len(calls) == 1 + assert slept == [] + + +def test_the_deferred_task_sends_every_run_s_scores_in_one_request(): + rec = _Recorder() + + def write_scores(ctx: Any) -> None: + for run in range(2): + ctx.score(f"trace-{run}", name="gate_passed", value=1.0, data_type="BOOLEAN") + ctx.score(f"trace-{run}", name="pass_at_k", value=1.0, data_type="BOOLEAN") + + with patch.object(_langfuse, "find_traces_per_conversation", return_value={}): + submit_trace_scoring( + lambda task, item_id="": task(), + _IDENTITY, + langfuse=rec.client, + dataset_item_id="item-1", + conversation_ids=["c0", "c1"], + window_start=_WINDOW[0], + window_end=_WINDOW[1], + suffix_runs=True, + write_scores=write_scores, + scope=ItemTraceScope("ds", True, "ds-1", "item-1"), + ) + + (bodies,) = _score_bodies(rec) + assert len(bodies) == 4 diff --git a/packages/gooddata-eval/tests/test_langfuse_client.py b/packages/gooddata-eval/tests/test_langfuse_client.py index 74ae5b8d0..de0e988d4 100644 --- a/packages/gooddata-eval/tests/test_langfuse_client.py +++ b/packages/gooddata-eval/tests/test_langfuse_client.py @@ -385,3 +385,104 @@ def test_a_dataset_run_item_for_an_unknown_item_raises(make_client): client = make_client(lambda request: httpx.Response(404, json={})) with pytest.raises(LookupError): client.api.dataset_run_items.create(run_name="run", dataset_item_id="local", trace_id="t-1") + + +def _score(name: str, value: float = 1.0, observation_id: str | None = None) -> dict: + return {"trace_id": "t-1", "name": name, "value": value, "data_type": "NUMERIC", "observation_id": observation_id} + + +def test_scores_are_posted_together_as_one_array(make_client): + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + return httpx.Response(202, json={"message": "Accepted"}) + + make_client(handler).create_scores([_score("quality_score", 0.5, "o-root"), _score("gate_passed")]) + + (request,) = seen + assert request.url.path == "/api/public/scores" + bodies = json.loads(request.content) + assert [(b["name"], b["value"], b.get("observationId")) for b in bodies] == [ + ("quality_score", 0.5, "o-root"), + ("gate_passed", 1.0, None), + ] + assert len({b["id"] for b in bodies}) == 2 + + +def test_an_unsendable_score_is_skipped_and_the_rest_of_the_batch_still_goes_out(make_client, caplog): + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + return httpx.Response(202, json={}) + + unknown_field = {**_score("bogus"), "unexpected": 1} + make_client(handler).create_scores([_score("a"), unknown_field, _score("nan", float("nan")), _score("b")]) + + (request,) = seen + assert [b["name"] for b in json.loads(request.content)] == ["a", "b"] + assert "bogus" in caplog.text + assert "nan" in caplog.text + + +def test_a_score_batch_larger_than_one_request_is_split(make_client, monkeypatch): + monkeypatch.setattr(client_module, "_SCORE_BATCH_SIZE", 2) + sizes: list[int] = [] + + def handler(request: httpx.Request) -> httpx.Response: + sizes.append(len(json.loads(request.content))) + return httpx.Response(202, json={}) + + make_client(handler).create_scores([_score(f"s{i}") for i in range(5)]) + + assert sizes == [2, 2, 1] + + +def test_a_throttled_batch_is_retried_whole_with_the_same_ids(make_client, monkeypatch): + slept: list[float] = [] + monkeypatch.setattr(client_module.time, "sleep", slept.append) + bodies: list[list[dict]] = [] + + def handler(request: httpx.Request) -> httpx.Response: + bodies.append(json.loads(request.content)) + if len(bodies) == 1: + return httpx.Response(429, headers={"Retry-After": "30"}, json={}) + return httpx.Response(202, json={}) + + make_client(handler).create_scores([_score("a"), _score("b")]) + + assert slept == [30.0] + assert [b["id"] for b in bodies[0]] == [b["id"] for b in bodies[1]] + + +def test_a_throttled_batch_is_not_retried_once_the_wait_is_refused(make_client): + waits: list[float] = [] + calls: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + calls.append(request) + return httpx.Response(429, headers={"Retry-After": "60"}, json={}) + + def refuse(delay: float) -> bool: + waits.append(delay) + return False + + with pytest.raises(httpx.HTTPStatusError): + make_client(handler).create_scores([_score("a")], wait=refuse) + + assert waits == [60.0] + assert len(calls) == 1 + + +def test_a_partly_rejected_batch_is_logged_and_not_retried(make_client, caplog): + calls: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + calls.append(request) + return httpx.Response(207, json={"accepted": 1, "rejected": 1, "errors": [{"index": 1, "message": "bad"}]}) + + make_client(handler).create_scores([_score("a"), _score("b")]) + + assert len(calls) == 1 + assert "1 of 2 scores rejected" in caplog.text diff --git a/packages/gooddata-eval/tests/test_langfuse_e2e_fake_server.py b/packages/gooddata-eval/tests/test_langfuse_e2e_fake_server.py index 76feccb9e..cd413dae5 100644 --- a/packages/gooddata-eval/tests/test_langfuse_e2e_fake_server.py +++ b/packages/gooddata-eval/tests/test_langfuse_e2e_fake_server.py @@ -119,7 +119,12 @@ def _attrs(span: dict) -> dict[str, Any]: def _score_bodies(server: FakeLangfuse) -> list[dict]: - return [call["json"] for call in server.calls("POST", _SCORES)] + """Every score posted, whether it went alone or inside an array.""" + bodies: list[dict] = [] + for call in server.calls("POST", _SCORES): + body = call["json"] + bodies.extend(body if isinstance(body, list) else [body]) + return bodies def test_agentic_inline_path_polls_looks_up_exports_and_scores(fake_langfuse: FakeLangfuse, capsys) -> None: @@ -349,12 +354,11 @@ def test_a_rate_limited_score_is_retried_and_lands(fake_langfuse: FakeLangfuse) with _agent_stubbed(["conv-1"]): _run_general_question(dataset_item_id="item-1", dataset_name="throttled") - bodies = _score_bodies(fake_langfuse) - # Fourteen writes -- seven scores on each of the gen-ai trace and the experiment span -- - # plus the one refused attempt the client repeated. - assert len(bodies) == 15 - posted_twice = [b for b in bodies if bodies.count(b) == 2] - assert len(posted_twice) == 2, "exactly one score body was posted twice" + # Fourteen scores -- seven on each of the gen-ai trace and the experiment span -- go out as + # one array; the refused request is repeated whole, with the same score ids. + (refused, landed) = [call["json"] for call in fake_langfuse.calls("POST", _SCORES)] + assert len(landed) == 14 + assert refused == landed def test_the_dataset_run_item_shim_exports_one_experiment_span(fake_langfuse: FakeLangfuse) -> None: From 45f68af00b651285af15c0b309ad918ea4142ed7 Mon Sep 17 00:00:00 2001 From: Jan Tychtl Date: Fri, 9 Oct 2026 13:03:36 +0200 Subject: [PATCH 5/5] feat(gooddata-eval): name gooddata-eval in Langfuse requests and spans Langfuse groups API traffic by User-Agent, so every gooddata-eval call showed up as python-httpx/0.28.1, indistinguishable from other callers. Its exported root spans carried only service.name, so a trace did not say which build or CI run produced it. Every client from make_http_client now sends gooddata-eval/ (; run=; sha=<12 chars of GITHUB_SHA>), with "local" for any value unset or unsafe in a header. The OTLP resource adds service.version and, when set, github.run_id, github.workflow, github.ref_name and github.event_name. jira: trivial risk: low --- .../src/gooddata_eval/core/langfuse/_env.py | 41 +++++++++- .../src/gooddata_eval/core/langfuse/otlp.py | 6 +- .../gooddata-eval/tests/test_langfuse_env.py | 76 ++++++++++++++++++- .../gooddata-eval/tests/test_langfuse_otlp.py | 3 +- 4 files changed, 120 insertions(+), 6 deletions(-) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/_env.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/_env.py index eee5e89dd..bf6ae0b3a 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/_env.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/_env.py @@ -5,10 +5,25 @@ import base64 import os +import re import httpx +from gooddata_eval._version import __version__ + _DEFAULT_BASE_URL = "https://us.cloud.langfuse.com" +SERVICE_NAME = "gooddata-eval" + +_LOCAL = "local" +_HEADER_SAFE = re.compile(r"[A-Za-z0-9._-]+") +_SHORT_SHA_LENGTH = 12 +# Resource attribute -> the GitHub Actions variable it is read from. +_GITHUB_RESOURCE = { + "github.run_id": "GITHUB_RUN_ID", + "github.workflow": "GITHUB_WORKFLOW", + "github.ref_name": "GITHUB_REF_NAME", + "github.event_name": "GITHUB_EVENT_NAME", +} def resolve_base_url() -> str: @@ -30,10 +45,34 @@ def basic_auth_header() -> str: return f"Basic {creds}" +def _header_value(name: str) -> str: + """The variable's value, or "local" when it is unset, blank or unsafe in a header.""" + value = os.environ.get(name, "").strip() + return value if _HEADER_SAFE.fullmatch(value) else _LOCAL + + +def user_agent() -> str: + """``gooddata-eval/ (; run=; sha=)``.""" + env = "ci" if os.environ.get("GITHUB_ACTIONS") == "true" else _LOCAL + sha = _header_value("GITHUB_SHA") + short_sha = sha[:_SHORT_SHA_LENGTH] if sha != _LOCAL else _LOCAL + return f"{SERVICE_NAME}/{__version__} ({env}; run={_header_value('GITHUB_RUN_ID')}; sha={short_sha})" + + +def resource_attributes() -> dict[str, str]: + """OTLP resource attributes of gooddata-eval's spans, with the GitHub run context when set.""" + attributes = {"service.name": SERVICE_NAME, "service.version": __version__} + for key, variable in _GITHUB_RESOURCE.items(): + value = os.environ.get(variable, "").strip() + if value: + attributes[key] = value + return attributes + + def make_http_client(*, timeout: float, transport: httpx.BaseTransport | None = None) -> httpx.Client: return httpx.Client( base_url=resolve_base_url(), - headers={"Authorization": basic_auth_header()}, + headers={"Authorization": basic_auth_header(), "User-Agent": user_agent()}, timeout=timeout, transport=transport, ) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/otlp.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/otlp.py index 2e605f734..efec8b373 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/otlp.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/otlp.py @@ -11,6 +11,7 @@ from typing import TYPE_CHECKING, Any from gooddata_eval._version import __version__ +from gooddata_eval.core.langfuse import _env if TYPE_CHECKING: import httpx @@ -42,7 +43,6 @@ ATTR_EXPERIMENT_ITEM_EXPECTED_OUTPUT = "langfuse.experiment.item.expected_output" ATTR_EXPERIMENT_ITEM_METADATA_PREFIX = "langfuse.experiment.item.metadata" -_SERVICE_NAME = "gooddata-eval" _EPOCH = datetime(1970, 1, 1, tzinfo=timezone.utc) @@ -112,7 +112,7 @@ class Span: def encode_export_request( - spans: list[Span], *, scope_name: str = _SERVICE_NAME, scope_version: str = __version__ + spans: list[Span], *, scope_name: str = _env.SERVICE_NAME, scope_version: str = __version__ ) -> dict[str, Any]: """Build the OTLP/JSON export request body for `POST /api/public/otel/v1/traces`.""" otlp_spans = [ @@ -131,7 +131,7 @@ def encode_export_request( return { "resourceSpans": [ { - "resource": {"attributes": [{"key": "service.name", "value": {"stringValue": _SERVICE_NAME}}]}, + "resource": {"attributes": [otlp_attribute(k, v) for k, v in _env.resource_attributes().items()]}, "scopeSpans": [{"scope": {"name": scope_name, "version": scope_version}, "spans": otlp_spans}], } ] diff --git a/packages/gooddata-eval/tests/test_langfuse_env.py b/packages/gooddata-eval/tests/test_langfuse_env.py index bd6ae2dc5..65301bd61 100644 --- a/packages/gooddata-eval/tests/test_langfuse_env.py +++ b/packages/gooddata-eval/tests/test_langfuse_env.py @@ -3,8 +3,38 @@ import base64 +import httpx import pytest -from gooddata_eval.core.langfuse._env import basic_auth_header, credentials_present, make_http_client, resolve_base_url +from gooddata_eval._version import __version__ +from gooddata_eval.core.langfuse._env import ( + basic_auth_header, + credentials_present, + make_http_client, + resolve_base_url, + resource_attributes, + user_agent, +) + +_GITHUB_ENV = { + "GITHUB_ACTIONS": "true", + "GITHUB_RUN_ID": "18273645", + "GITHUB_SHA": "0123456789abcdef0123456789abcdef01234567", + "GITHUB_WORKFLOW": "AI agent tests (staging)", + "GITHUB_REF_NAME": "master", + "GITHUB_EVENT_NAME": "schedule", +} + + +@pytest.fixture +def outside_ci(monkeypatch): + for name in _GITHUB_ENV: + monkeypatch.delenv(name, raising=False) + + +@pytest.fixture +def in_ci(monkeypatch): + for name, value in _GITHUB_ENV.items(): + monkeypatch.setenv(name, value) def test_resolve_base_url_prefers_base_url_over_host(monkeypatch): @@ -68,3 +98,47 @@ def test_make_http_client_uses_resolved_base_url_and_auth(monkeypatch): assert client.headers["Authorization"].startswith("Basic ") finally: client.close() + + +def test_a_ci_run_names_the_package_its_run_and_short_sha(in_ci): + assert user_agent() == f"gooddata-eval/{__version__} (ci; run=18273645; sha=0123456789ab)" + + +def test_outside_ci_the_user_agent_falls_back_to_local(outside_ci): + assert user_agent() == f"gooddata-eval/{__version__} (local; run=local; sha=local)" + + +@pytest.mark.parametrize("value", ["", "18273645\r\nX-Injected: 1", "run id"]) +def test_a_blank_or_unsafe_run_id_falls_back_to_local(outside_ci, monkeypatch, value): + monkeypatch.setenv("GITHUB_RUN_ID", value) + assert "run=local;" in user_agent() + + +def test_every_request_of_the_http_client_carries_the_user_agent(in_ci, monkeypatch): + monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk-test") + monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk-test") + sent: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + sent.append(request) + return httpx.Response(200, json={}) + + with make_http_client(timeout=5, transport=httpx.MockTransport(handler)) as client: + client.get("/api/public/v2/observations") + + assert sent[0].headers["User-Agent"] == user_agent() + + +def test_a_ci_run_puts_its_github_context_on_the_resource(in_ci): + assert resource_attributes() == { + "service.name": "gooddata-eval", + "service.version": __version__, + "github.run_id": "18273645", + "github.workflow": "AI agent tests (staging)", + "github.ref_name": "master", + "github.event_name": "schedule", + } + + +def test_outside_ci_the_resource_names_only_the_service(outside_ci): + assert resource_attributes() == {"service.name": "gooddata-eval", "service.version": __version__} diff --git a/packages/gooddata-eval/tests/test_langfuse_otlp.py b/packages/gooddata-eval/tests/test_langfuse_otlp.py index 0c6dd30f0..dbcef1cc0 100644 --- a/packages/gooddata-eval/tests/test_langfuse_otlp.py +++ b/packages/gooddata-eval/tests/test_langfuse_otlp.py @@ -7,6 +7,7 @@ import httpx import pytest +from gooddata_eval.core.langfuse._env import resource_attributes from gooddata_eval.core.langfuse.otlp import ( Span, encode_export_request, @@ -126,7 +127,7 @@ def test_encode_export_request_shape(): resource_spans = request["resourceSpans"] assert len(resource_spans) == 1 resource = resource_spans[0]["resource"] - assert resource["attributes"] == [{"key": "service.name", "value": {"stringValue": "gooddata-eval"}}] + assert {a["key"]: a["value"]["stringValue"] for a in resource["attributes"]} == resource_attributes() scope_spans = resource_spans[0]["scopeSpans"] assert len(scope_spans) == 1