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..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 @@ -10,7 +11,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 +23,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 +214,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 +272,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 +327,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 +362,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 +413,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 +483,89 @@ 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) + + +# 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. @@ -444,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) @@ -493,6 +632,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..5cb1f3988 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 _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,31 @@ 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, - ) + 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) 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, + ) ) - ) 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..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 @@ -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} @@ -579,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 @@ -587,7 +598,39 @@ 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 _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. + + 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/_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/client.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py index 1441457ec..d8cd75595 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,11 @@ from __future__ import annotations +import logging import threading import time import uuid +from collections.abc import Callable from datetime import datetime, timezone from typing import Any @@ -26,9 +28,16 @@ _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 +# 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 +47,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 +56,36 @@ 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) + + +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: @@ -142,26 +180,39 @@ 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. 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. + """ + for start in range(0, len(scores), _SCORE_BATCH_SIZE): + batch = [_score_body(**score) for score in scores[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.""" @@ -198,6 +249,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/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/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..313bf67e1 --- /dev/null +++ b/packages/gooddata-eval/tests/test_agentic_join_trace.py @@ -0,0 +1,335 @@ +# (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 +import threading +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, + 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 +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 + + +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_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_client.py b/packages/gooddata-eval/tests/test_langfuse_client.py index 931d7588e..5926c6e6b 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 @@ -377,3 +385,88 @@ 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_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: 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_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_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 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]