Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
183 changes: 179 additions & 4 deletions packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,14 +3,15 @@

from __future__ import annotations

import contextvars
import logging
import os
import threading
import time
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
Expand All @@ -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__)


Expand Down Expand Up @@ -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],
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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,
)


Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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.

Expand All @@ -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)
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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__)

Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -197,42 +215,39 @@ 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.

Deferred as a unit because ``build_run_context``'s workspace lookup and
``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)

Expand Down
Loading
Loading