From 89d8a892d36b6de55ba7dbf1f874de0abe45c7ba Mon Sep 17 00:00:00 2001 From: Jan Tychtl Date: Wed, 9 Sep 2026 18:26:19 +0200 Subject: [PATCH 1/3] feat(gooddata-eval): read Langfuse traces via v2 observations and move the client to core/langfuse Trace lookups query GET /api/public/v2/observations and fold the observation rows into one TraceSummary per trace: the root row carries session, metadata and latency, and cost is the sum over the trace's rows. HttpxLangfuseClient lives in core/langfuse/client.py, takes a transport for tests, and writes scores to POST /api/public/scores; core/agentic/_langfuse.py re-exports it and TraceSummary as _TraceObj for external importers. Dataset-run-item writes and the trace-version upsert stay on the legacy endpoints. jira: trivial risk: low --- .../gooddata_eval/core/agentic/_langfuse.py | 172 +-------- .../core/dataset/langfuse_source.py | 15 +- .../src/gooddata_eval/core/langfuse/client.py | 195 +++++++++++ .../core/langfuse/observations.py | 116 +++++++ .../tests/test_agentic_langfuse_trace.py | 96 +++--- .../tests/test_langfuse_client.py | 326 ++++++++++++++++++ .../tests/test_langfuse_observations.py | 234 +++++++++++++ 7 files changed, 918 insertions(+), 236 deletions(-) create mode 100644 packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py create mode 100644 packages/gooddata-eval/src/gooddata_eval/core/langfuse/observations.py create mode 100644 packages/gooddata-eval/tests/test_langfuse_client.py create mode 100644 packages/gooddata-eval/tests/test_langfuse_observations.py diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py index 5adef6d51..9394068df 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_langfuse.py @@ -3,12 +3,9 @@ from __future__ import annotations -import base64 import logging -import os import threading import time -import uuid from collections.abc import Iterator from contextlib import contextmanager from datetime import datetime, timedelta, timezone @@ -18,166 +15,13 @@ 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 -from gooddata_eval.core.langfuse._env import resolve_base_url +from gooddata_eval.core.langfuse._env import credentials_present +from gooddata_eval.core.langfuse.client import HttpxLangfuseClient -_log = logging.getLogger(__name__) - -# --------------------------------------------------------------------------- -# httpx-based Langfuse client — Python 3.14 safe (no Langfuse SDK required) -# --------------------------------------------------------------------------- +# 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 - -class _TraceObj: - """Duck-type wrapper around a raw Langfuse trace dict.""" - - def __init__(self, raw: dict) -> None: - self.id: str = raw.get("id", "") - self.metadata: dict = raw.get("metadata") or {} - self.session_id: str | None = raw.get("sessionId") or raw.get("session_id") - self.latency: float = float(raw.get("latency") or 0.0) - self.total_cost: float = float(raw.get("totalCost") or raw.get("total_cost") or 0.0) - - -class _TraceListResult: - def __init__(self, data: list[_TraceObj]) -> None: - self.data = data - - -class _TraceAPI: - def __init__(self, client: httpx.Client) -> None: - self._client = client - - def list( - self, from_timestamp: Any, to_timestamp: Any, limit: int, session_id: str | None = None - ) -> _TraceListResult: - """List traces in a window, optionally narrowed to one session server-side. - - ``session_id`` is what makes ``limit`` a non-issue. Without it the endpoint returns - every trace in the window newest-first, so an eval workspace busy enough to put more - than ``limit`` traces inside one item's window pushes that item's OWN (oldest) trace - off the page -- it then polls its whole retry budget against a page that can never - contain it, and the score orphans with only a generic "no trace found" line to show - for it. Concurrency makes that likelier by overlapping every item's window. Named - ``session_id`` because ``_fetch_traces_for_session`` probes for exactly that - parameter; gen-ai sets sessionId = conversationId. - """ - - def _ts(v: Any) -> str: - return v.isoformat() if hasattr(v, "isoformat") else str(v) - - params: dict[str, Any] = { - "fromTimestamp": _ts(from_timestamp), - "toTimestamp": _ts(to_timestamp), - "limit": limit, - } - # `is not None`, not truthiness: an empty id is a real filter value that matches - # nothing. Dropped here, the query would return the whole padded window for the - # caller's post-check to throw away, page after page, for the poll's whole budget. - if session_id is not None: - params["sessionId"] = session_id - resp = self._client.get("/api/public/traces", params=params) - resp.raise_for_status() - return _TraceListResult([_TraceObj(t) for t in resp.json().get("data", [])]) - - -class _DatasetRunItemsAPI: - def __init__(self, client: httpx.Client) -> None: - self._client = client - - def create( - self, - run_name: str, - dataset_item_id: str, - trace_id: str, - metadata: dict | None = None, - run_description: str = "", - ) -> None: - self._client.post( - "/api/public/dataset-run-items", - json={ - "runName": run_name, - "datasetItemId": dataset_item_id, - "traceId": trace_id, - "metadata": metadata or {}, - "runDescription": run_description, - }, - ).raise_for_status() - - -class _LangfuseAPI: - def __init__(self, client: httpx.Client) -> None: - self.trace = _TraceAPI(client) - self.dataset_run_items = _DatasetRunItemsAPI(client) - - -class HttpxLangfuseClient: - """Minimal Langfuse client using httpx — works on Python 3.14 (no Langfuse SDK needed).""" - - def __init__(self) -> None: - host = resolve_base_url() - pub = os.environ.get("LANGFUSE_PUBLIC_KEY", "") - sec = os.environ.get("LANGFUSE_SECRET_KEY", "") - if not pub or not sec: - raise RuntimeError( - "Langfuse credentials not set. " - "Export LANGFUSE_PUBLIC_KEY and LANGFUSE_SECRET_KEY before using --langfuse." - ) - creds = base64.b64encode(f"{pub}:{sec}".encode()).decode() - self._http = httpx.Client( - base_url=host, - headers={"Authorization": f"Basic {creds}"}, - timeout=10, - ) - self.api = _LangfuseAPI(self._http) - - def create_score( - self, - trace_id: str, - name: str, - value: float, - data_type: str, - comment: str | None = None, - ) -> None: - now = datetime.now(timezone.utc).isoformat() - # Langfuse API requires numeric value for BOOLEAN type (1.0/0.0), not JSON booleans - if isinstance(value, bool): - value = 1.0 if value else 0.0 - body: dict[str, Any] = { - "id": str(uuid.uuid4()), - "traceId": trace_id, - "name": name, - "value": value, - "dataType": data_type, - } - if comment: - body["comment"] = comment - self._http.post( - "/api/public/ingestion", - json={"batch": [{"id": str(uuid.uuid4()), "timestamp": now, "type": "score-create", "body": body}]}, - ).raise_for_status() - - def update_trace_version(self, trace_id: str, version: str) -> None: - """Upsert the trace version field via the ingestion endpoint.""" - now = datetime.now(timezone.utc).isoformat() - self._http.post( - "/api/public/ingestion", - json={ - "batch": [ - { - "id": str(uuid.uuid4()), - "timestamp": now, - "type": "trace-create", - "body": {"id": trace_id, "version": version}, - } - ] - }, - ).raise_for_status() - - def flush(self) -> None: - pass # no client-side batching - - def close(self) -> None: - self._http.close() +_log = logging.getLogger(__name__) def make_langfuse_client() -> HttpxLangfuseClient: @@ -191,7 +35,7 @@ def langfuse_credentials_present() -> bool: Separate from ``try_make_langfuse_client`` so a caller can ask the question without opening an httpx client it does not intend to use. """ - return bool(os.environ.get("LANGFUSE_PUBLIC_KEY")) and bool(os.environ.get("LANGFUSE_SECRET_KEY")) + return credentials_present() def try_make_langfuse_client() -> HttpxLangfuseClient | None: @@ -418,11 +262,7 @@ def _set_trace_version(langfuse: Any, trace_id: str, version: str) -> None: """Write model version into the Langfuse trace version field.""" try: if hasattr(langfuse, "update_trace_version"): - # HttpxLangfuseClient path langfuse.update_trace_version(trace_id, version) - elif hasattr(langfuse, "trace"): - # Langfuse Python SDK path (v2+) - langfuse.trace(id=trace_id, version=version) except Exception as exc: _log.warning("Failed to set trace version %r on %s: %s", version, trace_id, exc) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/dataset/langfuse_source.py b/packages/gooddata-eval/src/gooddata_eval/core/dataset/langfuse_source.py index 83c7d5f26..6c81ca179 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/dataset/langfuse_source.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/dataset/langfuse_source.py @@ -12,13 +12,11 @@ LANGFUSE_HOST — base URL, legacy alias for LANGFUSE_BASE_URL """ -import base64 -import os from typing import Any, TypeVar, cast import httpx -from gooddata_eval.core.langfuse._env import resolve_base_url +from gooddata_eval.core.langfuse._env import make_http_client from gooddata_eval.core.models import DatasetItem, SummaryInput _PAGE_SIZE = 100 @@ -28,16 +26,7 @@ def _make_client() -> httpx.Client: """Build an httpx client with Langfuse basic-auth headers.""" - host = resolve_base_url() - pub = os.environ.get("LANGFUSE_PUBLIC_KEY", "") - sec = os.environ.get("LANGFUSE_SECRET_KEY", "") - if not pub or not sec: - raise RuntimeError( - "Langfuse credentials not set. " - "Export LANGFUSE_PUBLIC_KEY and LANGFUSE_SECRET_KEY before using --langfuse-dataset." - ) - creds = base64.b64encode(f"{pub}:{sec}".encode()).decode() - return httpx.Client(base_url=host, headers={"Authorization": f"Basic {creds}"}, timeout=30) + return make_http_client(timeout=30) def _question_from_input(raw_input: Any) -> str: diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py new file mode 100644 index 000000000..a52402381 --- /dev/null +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py @@ -0,0 +1,195 @@ +# (C) 2026 GoodData Corporation +"""Minimal httpx Langfuse client: scores, OTLP export, dataset lookups and trace reads. + +No Langfuse SDK, so it works on every Python version the package supports. +""" + +from __future__ import annotations + +import threading +import time +import uuid +from datetime import datetime, timezone +from typing import Any + +import httpx + +from gooddata_eval.core.langfuse import _env, observations, otlp +from gooddata_eval.core.langfuse.observations import TraceSummary + +_SCORES_PATH = "/api/public/scores" +_OTLP_PATH = "/api/public/otel/v1/traces" +_INGESTION_PATH = "/api/public/ingestion" +_DATASET_RUN_ITEMS_PATH = "/api/public/dataset-run-items" + +_MAX_SCORE_ATTEMPTS = 3 +_DEFAULT_RETRY_DELAY = 0.5 +_MAX_RETRY_DELAY = 5.0 + + +def _is_retryable(resp: httpx.Response) -> bool: + return resp.status_code == 429 or resp.status_code >= 500 + + +def _retry_delay(resp: httpx.Response) -> float: + """Seconds to wait before the next attempt, from `Retry-After` when the server names one.""" + try: + return min(float(resp.headers.get("Retry-After", "")), _MAX_RETRY_DELAY) + except ValueError: + return _DEFAULT_RETRY_DELAY + + +class _TraceListResult: + def __init__(self, data: list[TraceSummary]) -> None: + self.data = data + + +class _TraceAPI: + """The `api.trace.list` shape external callers duck-type, served from observations.""" + + def __init__(self, owner: HttpxLangfuseClient) -> None: + self._owner = owner + + def list( + self, from_timestamp: Any, to_timestamp: Any, limit: int, session_id: str | None = None + ) -> _TraceListResult: + """List traces in a window, optionally narrowed to one session server-side. + + ``session_id`` is what makes ``limit`` a non-issue. Without it the window holds every + trace, newest-first, so an eval workspace busy enough to put more than ``limit`` + traces inside one item's window pushes that item's OWN (oldest) trace off the page -- + it then polls its whole retry budget against a page that can never contain it, and + the score orphans. Named ``session_id`` because ``_fetch_traces_for_session`` probes + for exactly that parameter; gen-ai sets sessionId = conversationId. + """ + return _TraceListResult( + self._owner.list_traces(from_time=from_timestamp, to_time=to_timestamp, limit=limit, session_id=session_id) + ) + + +class _DatasetRunItemsAPI: + def __init__(self, client: httpx.Client) -> None: + self._client = client + + def create( + self, + run_name: str, + dataset_item_id: str, + trace_id: str, + metadata: dict | None = None, + run_description: str = "", + ) -> None: + self._client.post( + _DATASET_RUN_ITEMS_PATH, + json={ + "runName": run_name, + "datasetItemId": dataset_item_id, + "traceId": trace_id, + "metadata": metadata or {}, + "runDescription": run_description, + }, + ).raise_for_status() + + +class _LangfuseAPI: + def __init__(self, owner: HttpxLangfuseClient) -> None: + self.trace = _TraceAPI(owner) + self.dataset_run_items = _DatasetRunItemsAPI(owner._http) + + +class HttpxLangfuseClient: + """Langfuse client over httpx, built from the standard Langfuse environment variables.""" + + def __init__(self, *, timeout: float = 10.0, transport: httpx.BaseTransport | None = None) -> None: + self._http = _env.make_http_client(timeout=timeout, transport=transport) + self._dataset_ids: dict[str, str | None] = {} + self._dataset_ids_lock = threading.Lock() + self.api = _LangfuseAPI(self) + + def create_score( + self, + trace_id: str, + name: str, + value: float, + data_type: str, + comment: str | None = None, + 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, + # BOOLEAN scores go over the wire as 1.0/0.0, not as JSON booleans. + "value": (1.0 if value else 0.0) if isinstance(value, bool) else value, + "dataType": data_type, + } + if comment: + body["comment"] = comment + if observation_id: + body["observationId"] = observation_id + resp = self._http.post(_SCORES_PATH, json=body) + for _retry in range(_MAX_SCORE_ATTEMPTS - 1): + if not _is_retryable(resp): + break + time.sleep(_retry_delay(resp)) + resp = self._http.post(_SCORES_PATH, json=body) + resp.raise_for_status() + + def export_spans(self, spans: list[otlp.Span]) -> None: + """Export spans to Langfuse over OTLP/HTTP JSON. Raises on a refused or rejected export.""" + resp = self._http.post( + _OTLP_PATH, + json=otlp.encode_export_request(spans), + headers={"x-langfuse-ingestion-version": "4"}, + ) + otlp.parse_export_response(resp) + + def dataset_id_for_item(self, item_id: str) -> str | None: + """The Langfuse dataset an item belongs to, or None when the id is not a Langfuse item. + + Cached because a run resolves the same handful of datasets once per item, from the + linking pool's worker threads. + """ + with self._dataset_ids_lock: + if item_id in self._dataset_ids: + return self._dataset_ids[item_id] + resp = self._http.get(f"/api/public/dataset-items/{item_id}") + if resp.status_code == 404: + dataset_id = None + else: + resp.raise_for_status() + dataset_id = resp.json().get("datasetId") + with self._dataset_ids_lock: + self._dataset_ids[item_id] = dataset_id + return dataset_id + + def list_traces( + self, *, from_time: Any, to_time: Any, limit: int, session_id: str | None = None + ) -> list[TraceSummary]: + return observations.list_traces_in_window( + self._http, from_time=from_time, to_time=to_time, limit=limit, session_id=session_id + ) + + def update_trace_version(self, trace_id: str, version: str) -> None: + """Upsert the trace version field via the ingestion endpoint.""" + now = datetime.now(timezone.utc).isoformat() + self._http.post( + _INGESTION_PATH, + json={ + "batch": [ + { + "id": str(uuid.uuid4()), + "timestamp": now, + "type": "trace-create", + "body": {"id": trace_id, "version": version}, + } + ] + }, + ).raise_for_status() + + def flush(self) -> None: + pass # no client-side batching + + def close(self) -> None: + self._http.close() diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/observations.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/observations.py new file mode 100644 index 000000000..8e2d5fbe6 --- /dev/null +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/observations.py @@ -0,0 +1,116 @@ +# (C) 2026 GoodData Corporation +"""Trace summaries read from the Langfuse v4 observations endpoint.""" + +from __future__ import annotations + +from datetime import datetime, timezone +from typing import TYPE_CHECKING, Any + +if TYPE_CHECKING: + import httpx + +_OBSERVATIONS_PATH = "/api/public/v2/observations" +# Everything a TraceSummary needs: core (ids, times, parent), basic (sessionId), usage +# (totalCost), metrics (latency), metadata. +_FIELDS = "core,basic,usage,metrics,metadata" + + +def _parse_time(value: Any) -> datetime | None: + """Parse an ISO 8601 timestamp into a timezone-aware datetime; anything else is None.""" + if not isinstance(value, str) or not value: + return None + try: + parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) + except ValueError: + return None + return parsed if parsed.tzinfo is not None else parsed.replace(tzinfo=timezone.utc) + + +def _iso(value: Any) -> str: + return value.isoformat() if hasattr(value, "isoformat") else str(value) + + +class TraceSummary: + """One Langfuse trace, seen through its root observation row or a legacy trace dict. + + An observation row identifies its trace by ``traceId`` and itself by ``id``; a legacy + trace dict carries the trace id in ``id`` and has no root observation. + """ + + def __init__(self, raw: dict, *, total_cost: float | None = None) -> None: + trace_id = raw.get("traceId") + self.id: str = trace_id or raw.get("id") or "" + self.root_observation_id: str | None = raw.get("id") if trace_id else None + self.metadata: dict = raw.get("metadata") or {} + self.session_id: str | None = raw.get("sessionId") or raw.get("session_id") + self.latency: float = float(raw.get("latency") or 0.0) + own_cost = raw.get("totalCost") or raw.get("total_cost") + self.total_cost: float = float((own_cost if total_cost is None else total_cost) or 0.0) + self.start_time: datetime | None = _parse_time(raw.get("startTime")) + self.end_time: datetime | None = _parse_time(raw.get("endTime")) + + +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. + """ + order: list[str] = [] + roots: dict[str, dict] = {} + costs: dict[str, float] = {} + 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] + + +def list_traces_in_window( + http: httpx.Client, + *, + from_time: Any, + to_time: Any, + limit: int, + session_id: str | None, + page_size: int = 500, + max_pages: int = 4, +) -> list[TraceSummary]: + """List up to ``limit`` traces whose observations start inside the window, newest first. + + ``session_id`` is sent whenever it is not None -- an empty id is a real filter value that + matches nothing, and dropping it would return the whole window for the caller to throw + away, page after page. + """ + params: dict[str, Any] = { + "fromStartTime": _iso(from_time), + "toStartTime": _iso(to_time), + "fields": _FIELDS, + "limit": page_size, + } + if session_id is not None: + params["sessionId"] = session_id + + rows: list[dict] = [] + summaries: list[TraceSummary] = [] + 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 []) + # Re-folded per page: a trace's root and its children can straddle a page boundary. + summaries = summarize_traces(rows) + cursor = (body.get("meta") or {}).get("cursor") + if not cursor or len(summaries) >= limit: + break + return summaries[:limit] diff --git a/packages/gooddata-eval/tests/test_agentic_langfuse_trace.py b/packages/gooddata-eval/tests/test_agentic_langfuse_trace.py index 56f252dca..770252870 100644 --- a/packages/gooddata-eval/tests/test_agentic_langfuse_trace.py +++ b/packages/gooddata-eval/tests/test_agentic_langfuse_trace.py @@ -1,6 +1,5 @@ # (C) 2026 GoodData Corporation. All rights reserved. # SPDX-License-Identifier: LicenseRef-GoodData-Enterprise -import os import time from datetime import datetime, timedelta, timezone from unittest.mock import MagicMock, patch @@ -13,9 +12,9 @@ _LINK_BUDGET_SEC, _MAX_DELAY, SKIP_ENV_VAR, + HttpxLangfuseClient, _fetch_traces_for_session, find_traces_per_conversation, - make_langfuse_client, observe, ) from gooddata_eval.core.agentic._trace_linker import ( @@ -266,20 +265,30 @@ def test_skip_switch_treats_explicit_off_values_as_off(monkeypatch, value, shoul # --- the session filter has to reach the server (M8) --- -def _stub_langfuse_http(monkeypatch, captured: list[dict]): - """A HttpxLangfuseClient whose httpx GET is recorded instead of sent.""" +def _observation_row(trace_id: str, session_id: str, latency: float) -> dict: + """A root observation row as /v2/observations returns it.""" + return { + "traceId": trace_id, + "id": f"o-{trace_id}", + "parentObservationId": None, + "sessionId": session_id, + "latency": latency, + "totalCost": 0.01, + } + + +def _stub_langfuse_http(monkeypatch, captured: list[httpx.Request], page: dict | None = None): + """A HttpxLangfuseClient whose requests are recorded instead of sent.""" monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk") monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk") monkeypatch.setenv("LANGFUSE_HOST", "https://lf.test") - client = make_langfuse_client() + body = page if page is not None else {"data": [], "meta": {}} - def _get(url, params=None, **_kw): - captured.append({"url": url, "params": params}) - return MagicMock(raise_for_status=lambda: None, json=lambda: {"data": []}) + def handler(request: httpx.Request) -> httpx.Response: + captured.append(request) + return httpx.Response(200, json=body) - client._http = MagicMock(get=_get) - client.api = type(client.api)(client._http) - return client + return HttpxLangfuseClient(transport=httpx.MockTransport(handler)) def test_the_trace_lookup_filters_by_session_server_side(monkeypatch): @@ -289,7 +298,7 @@ def test_the_trace_lookup_filters_by_session_server_side(monkeypatch): overlapping every item's window. The score then orphans after a full retry budget spent on a page that could never contain it. """ - captured: list[dict] = [] + captured: list[httpx.Request] = [] client = _stub_langfuse_http(monkeypatch, captured) _fetch_traces_for_session( @@ -297,32 +306,24 @@ def test_the_trace_lookup_filters_by_session_server_side(monkeypatch): ) assert len(captured) == 1 - assert captured[0]["params"]["sessionId"] == "conv-abc", "the session filter never reached the server" + assert captured[0].url.path == "/api/public/v2/observations" + assert dict(captured[0].url.params)["sessionId"] == "conv-abc", "the session filter never reached the server" def test_a_server_that_ignores_the_session_filter_cannot_hand_over_a_foreign_trace(monkeypatch): """The Langfuse API drops a query parameter it does not know rather than rejecting it. - The httpx client declares ``session_id``, which used to switch the local filter off for - it -- so a server that ignored the parameter returned the whole window, and the - max-latency pick attached this item's scores to a stranger's trace with no warning. - The server-side filter is still sent (paging); the local one is a post-check, not a - fallback. + A server that ignores the parameter answers with the whole window, and an unchecked + max-latency pick would then attach this item's scores to a stranger's trace with no + warning. The server-side filter is still sent (paging); the local one is a post-check, + not a fallback. """ - monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk") - monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk") - monkeypatch.setenv("LANGFUSE_HOST", "https://lf.test") - client = make_langfuse_client() + captured: list[httpx.Request] = [] page = { - "data": [ - {"id": "t-other", "sessionId": "conv-zzz", "latency": 9.0}, - {"id": "t-mine", "sessionId": "conv-abc", "latency": 1.0}, - ] + "data": [_observation_row("t-other", "conv-zzz", 9.0), _observation_row("t-mine", "conv-abc", 1.0)], + "meta": {}, } - client._http = MagicMock( - get=lambda url, params=None, **_kw: MagicMock(raise_for_status=lambda: None, json=lambda: page) - ) - client.api = type(client.api)(client._http) + client = _stub_langfuse_http(monkeypatch, captured, page) now = datetime.now(timezone.utc) found = _fetch_traces_for_session(client, "conv-abc", now, now, timedelta(seconds=2)) @@ -514,39 +515,20 @@ def _outer() -> None: assert seen == [False] -def test_an_empty_conversation_id_still_sends_the_server_side_filter(): +def test_an_empty_conversation_id_still_sends_the_server_side_filter(monkeypatch): """The filter must reach the server even when the id is empty. An empty id is a real filter value that matches nothing. Dropping the query parameter on a falsy id would fetch the entire padded window for the local post-check to throw away, so the poll would spend its whole budget on pages that can never match. """ - seen: list[dict] = [] - - def handler(request: httpx.Request) -> httpx.Response: - seen.append(dict(request.url.params)) - return httpx.Response(200, json={"data": []}) - - real_client = httpx.Client - - def fake_client(*args, **kwargs): - kwargs.pop("transport", None) - return real_client(*args, transport=httpx.MockTransport(handler), **kwargs) - - with ( - patch.dict( - os.environ, - {"LANGFUSE_HOST": "https://lf.test", "LANGFUSE_PUBLIC_KEY": "pk", "LANGFUSE_SECRET_KEY": "sk"}, - ), - patch.object(lf_module.httpx, "Client", fake_client), - ): - client = make_langfuse_client() + captured: list[httpx.Request] = [] + client = _stub_langfuse_http(monkeypatch, captured) now = datetime.now(timezone.utc) - with patch.object(lf_module.httpx, "Client", fake_client): - _fetch_traces_for_session(client, "", now - timedelta(minutes=5), now, timedelta(seconds=2)) + _fetch_traces_for_session(client, "", now - timedelta(minutes=5), now, timedelta(seconds=2)) - assert seen, "no request was made" - assert seen[0].get("sessionId") == "", ( - f"the empty id was dropped, so the server returned the whole window: {seen[0]}" - ) + assert captured, "no request was made" + assert captured[0].url.path == "/api/public/v2/observations" + params = dict(captured[0].url.params) + assert params.get("sessionId") == "", f"the empty id was dropped, so the server returned the whole window: {params}" diff --git a/packages/gooddata-eval/tests/test_langfuse_client.py b/packages/gooddata-eval/tests/test_langfuse_client.py new file mode 100644 index 000000000..49a59daa9 --- /dev/null +++ b/packages/gooddata-eval/tests/test_langfuse_client.py @@ -0,0 +1,326 @@ +# (C) 2026 GoodData Corporation +from __future__ import annotations + +import base64 +import json +import re +from datetime import datetime, timedelta, timezone + +import httpx +import pytest +from gooddata_eval.core.langfuse import client as client_module +from gooddata_eval.core.langfuse.client import HttpxLangfuseClient +from gooddata_eval.core.langfuse.observations import TraceSummary +from gooddata_eval.core.langfuse.otlp import Span, otlp_attribute + +_BASIC_AUTH = f"Basic {base64.b64encode(b'pk:sk').decode()}" + + +@pytest.fixture(autouse=True) +def _langfuse_env(monkeypatch): + monkeypatch.setenv("LANGFUSE_BASE_URL", "https://lf.test") + monkeypatch.delenv("LANGFUSE_HOST", raising=False) + monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk") + monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk") + + +@pytest.fixture +def make_client(): + """Build clients against a MockTransport handler and close them at teardown.""" + created: list[HttpxLangfuseClient] = [] + + def _make(handler) -> HttpxLangfuseClient: + client = HttpxLangfuseClient(transport=httpx.MockTransport(handler)) + created.append(client) + return client + + yield _make + for client in created: + client.close() + + +def _ok(request: httpx.Request) -> httpx.Response: + return httpx.Response(200, json={}) + + +def _span() -> Span: + start = datetime(2026, 9, 9, 10, 0, tzinfo=timezone.utc) + return Span( + trace_id="0" * 32, + span_id="1" * 16, + name="gd-eval", + start=start, + end=start + timedelta(seconds=2), + attributes=[otlp_attribute("langfuse.observation.type", "span")], + ) + + +def test_missing_credentials_are_refused(monkeypatch): + monkeypatch.delenv("LANGFUSE_PUBLIC_KEY", raising=False) + monkeypatch.delenv("LANGFUSE_SECRET_KEY", raising=False) + with pytest.raises(RuntimeError, match="credentials"): + HttpxLangfuseClient() + + +def test_a_score_is_posted_to_the_scores_endpoint(make_client): + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + return _ok(request) + + make_client(handler).create_score("t-1", "quality_score", 0.75, "NUMERIC", comment="3/4 checks passed") + + assert seen[0].method == "POST" + assert seen[0].url.path == "/api/public/scores" + assert seen[0].headers["Authorization"] == _BASIC_AUTH + body = json.loads(seen[0].content) + assert body["traceId"] == "t-1" + assert body["name"] == "quality_score" + assert body["value"] == 0.75 + assert body["dataType"] == "NUMERIC" + assert body["comment"] == "3/4 checks passed" + assert re.fullmatch(r"[0-9a-f-]{36}", body["id"]) + assert "observationId" not in body + + +def test_the_observation_id_is_sent_only_when_asked_for(make_client): + bodies: list[dict] = [] + + def handler(request: httpx.Request) -> httpx.Response: + bodies.append(json.loads(request.content)) + return _ok(request) + + client = make_client(handler) + client.create_score("t-1", "pass_at_k", 1.0, "BOOLEAN") + client.create_score("t-1", "pass_at_k", 1.0, "BOOLEAN", observation_id="o-root") + + assert "observationId" not in bodies[0] + assert bodies[1]["observationId"] == "o-root" + + +def test_a_boolean_score_is_sent_as_a_number(make_client): + bodies: list[dict] = [] + + def handler(request: httpx.Request) -> httpx.Response: + bodies.append(json.loads(request.content)) + return _ok(request) + + client = make_client(handler) + client.create_score("t-1", "pass_at_k", True, "BOOLEAN") + client.create_score("t-1", "pass_at_k", False, "BOOLEAN") + + assert [body["value"] for body in bodies] == [1.0, 0.0] + + +def test_a_throttled_score_is_retried_after_the_delay_the_server_asked_for(make_client, monkeypatch): + slept: list[float] = [] + monkeypatch.setattr(client_module.time, "sleep", slept.append) + statuses = [429, 200] + + def handler(request: httpx.Request) -> httpx.Response: + status = statuses[len(slept)] + return httpx.Response(status, headers={"Retry-After": "2"} if status == 429 else {}, json={}) + + make_client(handler).create_score("t-1", "quality_score", 1.0, "NUMERIC") + + assert slept == [2.0] + + +def test_a_retry_delay_the_server_did_not_name_falls_back_to_a_short_one(make_client, monkeypatch): + slept: list[float] = [] + monkeypatch.setattr(client_module.time, "sleep", slept.append) + statuses = [500, 200] + + make_client(lambda request: httpx.Response(statuses[len(slept)], json={})).create_score( + "t-1", "quality_score", 1.0, "NUMERIC" + ) + + assert slept == [0.5] + + +def test_a_score_gives_up_after_three_attempts(make_client, monkeypatch): + monkeypatch.setattr(client_module.time, "sleep", lambda _seconds: None) + attempts = 0 + + def handler(request: httpx.Request) -> httpx.Response: + nonlocal attempts + attempts += 1 + return httpx.Response(503, json={}) + + client = make_client(handler) + with pytest.raises(httpx.HTTPStatusError): + client.create_score("t-1", "quality_score", 1.0, "NUMERIC") + + assert attempts == 3 + + +def test_a_rejected_score_is_not_retried(make_client): + attempts = 0 + + def handler(request: httpx.Request) -> httpx.Response: + nonlocal attempts + attempts += 1 + return httpx.Response(400, json={}) + + client = make_client(handler) + with pytest.raises(httpx.HTTPStatusError): + client.create_score("t-1", "quality_score", 1.0, "NUMERIC") + + assert attempts == 1 + + +def test_spans_are_exported_as_otlp_with_the_v4_ingestion_header(make_client): + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + return httpx.Response(200, json={"partialSuccess": {}}) + + make_client(handler).export_spans([_span()]) + + assert seen[0].method == "POST" + assert seen[0].url.path == "/api/public/otel/v1/traces" + assert seen[0].headers["x-langfuse-ingestion-version"] == "4" + assert seen[0].headers["Authorization"] == _BASIC_AUTH + span = json.loads(seen[0].content)["resourceSpans"][0]["scopeSpans"][0]["spans"][0] + assert re.fullmatch(r"[0-9a-f]{32}", span["traceId"]) + assert re.fullmatch(r"[0-9a-f]{16}", span["spanId"]) + assert int(span["startTimeUnixNano"]) <= int(span["endTimeUnixNano"]) + + +def test_a_refused_export_raises(make_client): + client = make_client(lambda request: httpx.Response(400, text="bad span")) + with pytest.raises(RuntimeError): + client.export_spans([_span()]) + + +def test_a_partially_rejected_export_raises(make_client): + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response(200, json={"partialSuccess": {"rejectedSpans": 1, "errorMessage": "nope"}}) + + client = make_client(handler) + with pytest.raises(RuntimeError, match="nope"): + client.export_spans([_span()]) + + +def test_the_dataset_of_an_item_is_looked_up_once_and_cached(make_client): + calls: list[str] = [] + + def handler(request: httpx.Request) -> httpx.Response: + calls.append(request.url.path) + return httpx.Response(200, json={"id": "item-1", "datasetId": "ds-1"}) + + client = make_client(handler) + assert client.dataset_id_for_item("item-1") == "ds-1" + assert client.dataset_id_for_item("item-1") == "ds-1" + assert calls == ["/api/public/dataset-items/item-1"] + + +def test_an_unknown_dataset_item_is_none_and_is_not_asked_for_twice(make_client): + calls = 0 + + def handler(request: httpx.Request) -> httpx.Response: + nonlocal calls + calls += 1 + return httpx.Response(404, json={"message": "not found"}) + + client = make_client(handler) + assert client.dataset_id_for_item("local-item") is None + assert client.dataset_id_for_item("local-item") is None + assert calls == 1 + + +def test_a_broken_dataset_item_lookup_raises(make_client): + client = make_client(lambda request: httpx.Response(500, text="boom")) + with pytest.raises(httpx.HTTPStatusError): + client.dataset_id_for_item("item-1") + + +def test_the_compat_trace_api_returns_trace_summaries(make_client): + root = { + "traceId": "t-1", + "id": "o-root", + "parentObservationId": None, + "sessionId": "conv-1", + "latency": 4.0, + "totalCost": None, + } + child = {"traceId": "t-1", "id": "o-child", "parentObservationId": "o-root", "totalCost": 0.02} + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + return httpx.Response(200, json={"data": [root, child], "meta": {}}) + + now = datetime.now(timezone.utc) + result = make_client(handler).api.trace.list(from_timestamp=now, to_timestamp=now, limit=100, session_id="conv-1") + + assert seen[0].url.path == "/api/public/v2/observations" + assert dict(seen[0].url.params)["sessionId"] == "conv-1" + assert all(isinstance(trace, TraceSummary) for trace in result.data) + assert [(trace.id, trace.session_id, trace.latency, trace.total_cost) for trace in result.data] == [ + ("t-1", "conv-1", 4.0, 0.02) + ] + + +def test_the_compat_trace_api_accepts_timestamp_strings(make_client): + seen: list[dict] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(dict(request.url.params)) + return httpx.Response(200, json={"data": [], "meta": {}}) + + make_client(handler).api.trace.list( + from_timestamp="2026-09-09T10:00:00+00:00", to_timestamp="2026-09-09T10:05:00+00:00", limit=10 + ) + + assert seen[0]["fromStartTime"] == "2026-09-09T10:00:00+00:00" + assert "sessionId" not in seen[0] + + +def test_a_dataset_run_item_is_posted_to_the_legacy_endpoint(make_client): + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + return _ok(request) + + make_client(handler).api.dataset_run_items.create( + run_name="ds_2026_model", + dataset_item_id="item-1", + trace_id="t-1", + metadata={"model_version": "m"}, + run_description="desc", + ) + + assert seen[0].method == "POST" + assert seen[0].url.path == "/api/public/dataset-run-items" + assert json.loads(seen[0].content) == { + "runName": "ds_2026_model", + "datasetItemId": "item-1", + "traceId": "t-1", + "metadata": {"model_version": "m"}, + "runDescription": "desc", + } + + +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(httpx.HTTPStatusError): + client.api.dataset_run_items.create(run_name="run", dataset_item_id="local", trace_id="t-1") + + +def test_the_trace_version_upsert_uses_the_ingestion_endpoint(make_client): + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + return _ok(request) + + make_client(handler).update_trace_version("t-1", "gpt-5.2") + + assert seen[0].url.path == "/api/public/ingestion" + event = json.loads(seen[0].content)["batch"][0] + assert event["type"] == "trace-create" + assert event["body"] == {"id": "t-1", "version": "gpt-5.2"} diff --git a/packages/gooddata-eval/tests/test_langfuse_observations.py b/packages/gooddata-eval/tests/test_langfuse_observations.py new file mode 100644 index 000000000..b942e64e4 --- /dev/null +++ b/packages/gooddata-eval/tests/test_langfuse_observations.py @@ -0,0 +1,234 @@ +# (C) 2026 GoodData Corporation +from __future__ import annotations + +from datetime import datetime, timezone + +import httpx +import pytest +from gooddata_eval.core.langfuse.observations import TraceSummary, list_traces_in_window, summarize_traces + + +def _row( + trace_id: str, + obs_id: str, + *, + parent: str | None = None, + latency: float | None = None, + total_cost: float | None = None, + session_id: str | None = None, + metadata: dict | None = None, +) -> dict: + return { + "traceId": trace_id, + "id": obs_id, + "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", + } + + +def test_rows_are_grouped_into_one_summary_per_trace(): + rows = [ + _row("t-1", "o-root", latency=3.0), + _row("t-1", "o-child", parent="o-root", total_cost=0.01), + _row("t-2", "o-root-2", latency=1.0), + ] + assert [summary.id for summary in summarize_traces(rows)] == ["t-1", "t-2"] + + +def test_the_root_is_the_row_without_a_parent_observation(): + rows = [ + _row("t-1", "o-child", parent="o-root", session_id="not-the-root", latency=99.0), + _row("t-1", "o-root", latency=3.0, session_id="conv-1", metadata={"conversation_id": "conv-1"}), + ] + (summary,) = summarize_traces(rows) + assert summary.root_observation_id == "o-root" + assert summary.session_id == "conv-1" + assert summary.metadata == {"conversation_id": "conv-1"} + assert summary.latency == 3.0 + + +def test_the_trace_cost_is_the_sum_over_its_rows(): + # The gen-ai root row carries no cost of its own; the model calls under it do. + rows = [ + _row("t-1", "o-root", latency=3.0, total_cost=None), + _row("t-1", "o-a", parent="o-root", total_cost=0.012), + _row("t-1", "o-b", parent="o-root", total_cost=0.008), + ] + (summary,) = summarize_traces(rows) + assert summary.total_cost == pytest.approx(0.02) + + +def test_a_trace_whose_root_is_not_on_the_page_is_dropped(): + # A poll that catches a conversation mid-ingestion sees children only; such a trace has + # no latency and no session of its own, so it must not be offered to the caller. + rows = [_row("t-partial", "o-child", parent="o-missing", total_cost=0.01)] + assert summarize_traces(rows) == [] + + +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) + assert summary.end_time == datetime(2026, 9, 9, 10, 0, 12, tzinfo=timezone.utc) + + +def test_a_legacy_trace_dict_is_still_accepted(): + summary = TraceSummary( + {"id": "t-1", "sessionId": "conv-1", "latency": 4.5, "totalCost": 0.03, "metadata": {"k": "v"}} + ) + assert (summary.id, summary.session_id, summary.latency, summary.total_cost) == ("t-1", "conv-1", 4.5, 0.03) + assert summary.metadata == {"k": "v"} + assert summary.root_observation_id is None + assert summary.start_time is None + + +def test_missing_numbers_read_as_zero(): + summary = TraceSummary({"id": "t-1"}) + assert summary.latency == 0.0 + assert summary.total_cost == 0.0 + assert summary.metadata == {} + + +def _client(handler) -> httpx.Client: + return httpx.Client(base_url="https://lf.test", transport=httpx.MockTransport(handler)) + + +def test_the_window_query_carries_every_required_parameter(): + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + return httpx.Response(200, json={"data": [_row("t-1", "o-root", latency=1.0)], "meta": {}}) + + with _client(handler) as http: + found = list_traces_in_window( + http, + from_time=datetime(2026, 9, 9, 10, 0, tzinfo=timezone.utc), + to_time=datetime(2026, 9, 9, 10, 5, tzinfo=timezone.utc), + limit=100, + session_id="conv-1", + page_size=250, + ) + + assert [summary.id for summary in found] == ["t-1"] + assert seen[0].url.path == "/api/public/v2/observations" + params = dict(seen[0].url.params) + assert params["fromStartTime"] == "2026-09-09T10:00:00+00:00" + assert params["toStartTime"] == "2026-09-09T10:05:00+00:00" + assert params["sessionId"] == "conv-1" + assert params["fields"] == "core,basic,usage,metrics,metadata" + assert params["limit"] == "250" + + +def test_an_empty_session_id_is_still_sent(): + # An empty id is a real filter value that matches nothing; dropped, the query returns + # the whole window for the caller to throw away. + seen: list[dict] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(dict(request.url.params)) + return httpx.Response(200, json={"data": [], "meta": {}}) + + now = datetime.now(timezone.utc) + with _client(handler) as http: + list_traces_in_window(http, from_time=now, to_time=now, limit=10, session_id="") + + assert seen[0]["sessionId"] == "" + + +def test_no_session_filter_is_sent_when_none_is_asked_for(): + seen: list[dict] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(dict(request.url.params)) + return httpx.Response(200, json={"data": [], "meta": {}}) + + now = datetime.now(timezone.utc) + with _client(handler) as http: + list_traces_in_window(http, from_time=now, to_time=now, limit=10, session_id=None) + + assert "sessionId" not in seen[0] + + +def test_the_cursor_is_followed_until_the_server_stops_handing_one_out(): + pages = [ + {"data": [_row("t-1", "o-1", parent="o-root-1", total_cost=0.5)], "meta": {"cursor": "c1"}}, + {"data": [_row("t-1", "o-root-1", latency=2.0)], "meta": {}}, + ] + cursors: list[str | None] = [] + + def handler(request: httpx.Request) -> httpx.Response: + cursors.append(request.url.params.get("cursor")) + return httpx.Response(200, json=pages[len(cursors) - 1]) + + now = datetime.now(timezone.utc) + with _client(handler) as http: + found = list_traces_in_window(http, from_time=now, to_time=now, limit=10, session_id=None) + + # The root arrives on the second page, so the trace is only complete once both are in. + assert cursors == [None, "c1"] + assert [(summary.id, summary.total_cost, summary.latency) for summary in found] == [("t-1", 0.5, 2.0)] + + +def test_paging_stops_at_max_pages(): + calls = 0 + + def handler(request: httpx.Request) -> httpx.Response: + nonlocal calls + calls += 1 + return httpx.Response(200, json={"data": [], "meta": {"cursor": "always-more"}}) + + now = datetime.now(timezone.utc) + with _client(handler) as http: + list_traces_in_window(http, from_time=now, to_time=now, limit=10, session_id=None, max_pages=3) + + assert calls == 3 + + +def test_paging_stops_once_enough_traces_are_collected(): + calls = 0 + + def handler(request: httpx.Request) -> httpx.Response: + nonlocal calls + calls += 1 + return httpx.Response( + 200, + json={ + "data": [_row(f"t-{calls}-a", "o-a", latency=1.0), _row(f"t-{calls}-b", "o-b", latency=1.0)], + "meta": {"cursor": "more"}, + }, + ) + + now = datetime.now(timezone.utc) + with _client(handler) as http: + found = list_traces_in_window(http, from_time=now, to_time=now, limit=2, session_id=None) + + assert calls == 1 + assert len(found) == 2 + + +def test_more_traces_than_asked_for_are_truncated(): + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response( + 200, + json={"data": [_row(f"t-{i}", f"o-{i}", latency=1.0) for i in range(5)], "meta": {}}, + ) + + now = datetime.now(timezone.utc) + with _client(handler) as http: + found = list_traces_in_window(http, from_time=now, to_time=now, limit=3, session_id=None) + + assert [summary.id for summary in found] == ["t-0", "t-1", "t-2"] + + +def test_a_failed_page_raises(): + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response(500, text="boom") + + 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) From bbe11871bb42ce1425d82226e0217e914bd7c8bb Mon Sep 17 00:00:00 2001 From: Jan Tychtl Date: Wed, 9 Sep 2026 18:41:33 +0200 Subject: [PATCH 2/3] test(gooddata-eval): add an in-process fake Langfuse server for wire-level tests Adds tests/_fake_langfuse.py, a threading.Thread-hosted http.server that answers the five Langfuse v4 endpoints the package uses plus the three legacy ones, records every request, and is configurable per test (dataset items, missing ids, observation paging/empty-first-call, OTLP status/body, scores 429-once). It doubles as a pytest fixture (fake_langfuse in conftest.py, which points the real Langfuse env vars at it) and as a runnable script that prints and records the wire traffic for manual inspection. tests/test_fake_langfuse.py exercises the fixture itself with a real httpx.Client, no mocks. jira: trivial risk: low --- .../gooddata-eval/tests/_fake_langfuse.py | 294 ++++++++++++++++++ packages/gooddata-eval/tests/conftest.py | 13 + .../gooddata-eval/tests/test_fake_langfuse.py | 140 +++++++++ 3 files changed, 447 insertions(+) create mode 100644 packages/gooddata-eval/tests/_fake_langfuse.py create mode 100644 packages/gooddata-eval/tests/test_fake_langfuse.py diff --git a/packages/gooddata-eval/tests/_fake_langfuse.py b/packages/gooddata-eval/tests/_fake_langfuse.py new file mode 100644 index 000000000..df342cc62 --- /dev/null +++ b/packages/gooddata-eval/tests/_fake_langfuse.py @@ -0,0 +1,294 @@ +# (C) 2026 GoodData Corporation +"""In-process fake Langfuse HTTP server: a pytest fixture and a runnable wire-watching script. + +A `threading.Thread`-hosted `http.server` answering the Langfuse v4 endpoints the package +uses, plus the three legacy ones, with canned/synthesised data recorded on `requests`. +""" + +from __future__ import annotations + +import argparse +import hashlib +import itertools +import json +import signal +import sys +from datetime import datetime, timedelta, timezone +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path +from threading import Thread +from urllib.parse import parse_qs, urlsplit + +_DATASET_ITEMS_PATH = "/api/public/dataset-items" +_OBSERVATIONS_PATH = "/api/public/v2/observations" +_OTLP_PATH = "/api/public/otel/v1/traces" +_SCORES_PATH = "/api/public/scores" +_TRACES_PATH = "/api/public/traces" +_INGESTION_PATH = "/api/public/ingestion" +_DATASET_RUN_ITEMS_PATH = "/api/public/dataset-run-items" + +_ROOT_LATENCY_SECONDS = 12.5 +_CHILD_COSTS = (0.01, 0.02) + + +def _short_hash(value: str, length: int) -> str: + return hashlib.md5(value.encode()).hexdigest()[:length] + + +def observation_rows(session_id: str) -> list[dict]: + """A gen-ai-shaped root row plus two child rows for one conversation session.""" + trace_id = _short_hash(session_id, 32) + root_id = _short_hash(f"{session_id}:root", 16) + start = datetime(2026, 1, 1, tzinfo=timezone.utc) + times = {"startTime": start.isoformat(), "endTime": (start + timedelta(seconds=_ROOT_LATENCY_SECONDS)).isoformat()} + shared = {"traceId": trace_id, "type": "GENERATION", "sessionId": session_id, **times} + root = { + "id": root_id, + "name": "conversation.send_message", + "parentObservationId": None, + "isRootObservation": True, + "latency": _ROOT_LATENCY_SECONDS, + "totalCost": None, + "metadata": {"conversation_id": session_id}, + **shared, + } + children = [ + { + "id": _short_hash(f"{session_id}:child{i}", 16), + "name": f"model.call.{i}", + "parentObservationId": root_id, + "isRootObservation": False, + "latency": None, + "totalCost": cost, + "metadata": {}, + **shared, + } + for i, cost in enumerate(_CHILD_COSTS, start=1) + ] + return [root, *children] + + +def _canned_item(item_id: str, dataset_name: str) -> dict: + return { + "id": item_id, + "datasetName": dataset_name, + "input": {"question": f"Canned question for {item_id}"}, + "expectedOutput": f"Rubric: the answer for {item_id} must directly address the question.", + "metadata": {}, + } + + +class FakeLangfuse: + """A fake Langfuse server on `127.0.0.1:0`, configurable per test before each request.""" + + def __init__(self, *, port: int = 0, verbose: bool = False) -> None: + self.dataset_name = "fake" + self.dataset_id = "ds-fake" + self.items: list[dict] = [] + self.missing: set[str] = set() + self.first_observations_call_empty = False + self.observations_pages = 1 + self.otlp_status = 200 + self.otlp_body: dict = {} + self.scores_429_once = False + self.ingestion_body: dict = {"successes": [], "errors": []} + self.verbose = verbose + self.requests: list[dict] = [] + self.on_request = None + handler = self._make_handler() + self._httpd = ThreadingHTTPServer(("127.0.0.1", port), handler) + self._thread = Thread(target=self._httpd.serve_forever, daemon=True) + + @property + def base_url(self) -> str: + host, port = self._httpd.server_address[:2] + return f"http://{host}:{port}" + + def start(self) -> None: + self._thread.start() + + def stop(self) -> None: + self._httpd.shutdown() + self._httpd.server_close() + + def __enter__(self) -> FakeLangfuse: + self.start() + return self + + def __exit__(self, *exc_info: object) -> None: + self.stop() + + def calls(self, method: str, path_prefix: str) -> list[dict]: + return [r for r in self.requests if r["method"] == method and r["path"].startswith(path_prefix)] + + def _make_handler(self) -> type[BaseHTTPRequestHandler]: + server = self + + class Handler(BaseHTTPRequestHandler): + def log_message(self, format: str, *args: object) -> None: + if server.verbose: + super().log_message(format, *args) + + def do_GET(self) -> None: + server._handle(self, "GET") + + def do_POST(self) -> None: + server._handle(self, "POST") + + return Handler + + def _handle(self, handler: BaseHTTPRequestHandler, method: str) -> None: + try: + self._route(handler, method) + except Exception as exc: # never let a bad request kill the server thread + self._respond(handler, 400, {"message": f"fake_langfuse error: {exc}"}) + + def _route(self, handler: BaseHTTPRequestHandler, method: str) -> None: + parts = urlsplit(handler.path) + path = parts.path + query = {k: v[0] for k, v in parse_qs(parts.query).items()} + body = self._read_json_body(handler) + record = { + "method": method, + "path": path, + "query": query, + "headers": {k.lower(): v for k, v in handler.headers.items()}, + "json": body, + } + self.requests.append(record) + if self.on_request is not None: + self.on_request(record) + + if method == "GET" and path == _DATASET_ITEMS_PATH: + self._get_dataset_items(handler, query) + elif method == "GET" and path.startswith(f"{_DATASET_ITEMS_PATH}/"): + self._get_dataset_item(handler, path.rsplit("/", 1)[-1]) + elif method == "GET" and path == _OBSERVATIONS_PATH: + self._get_observations(handler, query) + elif method == "POST" and path == _OTLP_PATH: + self._respond(handler, self.otlp_status, self.otlp_body) + elif method == "POST" and path == _SCORES_PATH: + self._post_scores(handler) + elif method == "GET" and path == _TRACES_PATH: + self._get_traces(handler, query) + elif method == "POST" and path == _INGESTION_PATH: + self._respond(handler, 200, self.ingestion_body) + elif method == "POST" and path == _DATASET_RUN_ITEMS_PATH: + self._respond(handler, 200, {}) + else: + self._respond(handler, 404, {"message": f"fake_langfuse: no route for {method} {path}"}) + + @staticmethod + def _read_json_body(handler: BaseHTTPRequestHandler) -> dict | None: + length = int(handler.headers.get("Content-Length") or 0) + if not length: + return None + raw = handler.rfile.read(length) + return json.loads(raw) if raw else None + + def _get_dataset_items(self, handler: BaseHTTPRequestHandler, query: dict) -> None: + if query.get("datasetName") != self.dataset_name: + self._respond(handler, 404, {"message": "Dataset not found"}) + return + meta = {"page": 1, "limit": 100, "totalItems": len(self.items), "totalPages": 1} + self._respond(handler, 200, {"data": self.items, "meta": meta}) + + def _get_dataset_item(self, handler: BaseHTTPRequestHandler, item_id: str) -> None: + if item_id in self.missing: + self._respond(handler, 404, {"message": "Dataset item not found"}) + return + self._respond(handler, 200, {"id": item_id, "datasetId": self.dataset_id, "datasetName": self.dataset_name}) + + def _post_scores(self, handler: BaseHTTPRequestHandler) -> None: + if self.scores_429_once: + self.scores_429_once = False + self._respond(handler, 429, {"message": "rate limited"}, headers={"Retry-After": "0"}) + return + self._respond(handler, 200, {}) + + def _get_observations(self, handler: BaseHTTPRequestHandler, query: dict) -> None: + session_id = query.get("sessionId", "") + if self.first_observations_call_empty: + self.first_observations_call_empty = False + self._respond(handler, 200, {"data": [], "meta": {}}) + return + rows = observation_rows(session_id) + if self.observations_pages == 2 and query.get("cursor") is None: + self._respond(handler, 200, {"data": rows[:1], "meta": {"cursor": "page-2"}}) + return + if self.observations_pages == 2: + self._respond(handler, 200, {"data": rows[1:], "meta": {}}) + return + self._respond(handler, 200, {"data": rows, "meta": {}}) + + def _get_traces(self, handler: BaseHTTPRequestHandler, query: dict) -> None: + session_id = query.get("sessionId", "") + legacy_trace = { + "id": _short_hash(session_id, 32), + "sessionId": session_id, + "latency": _ROOT_LATENCY_SECONDS, + "totalCost": sum(_CHILD_COSTS), + "metadata": {"conversation_id": session_id}, + } + self._respond(handler, 200, {"data": [legacy_trace]}) + + def _respond( + self, handler: BaseHTTPRequestHandler, status: int, body: dict, *, headers: dict | None = None + ) -> None: + payload = json.dumps(body).encode() + handler.send_response(status) + handler.send_header("Content-Type", "application/json") + handler.send_header("Content-Length", str(len(payload))) + for key, value in (headers or {}).items(): + handler.send_header(key, value) + handler.end_headers() + handler.wfile.write(payload) + + +def _parse_args(argv: list[str]) -> argparse.Namespace: + parser = argparse.ArgumentParser(description="Run a fake Langfuse server and watch the wire.") + parser.add_argument("--port", type=int, default=8765) + parser.add_argument("--dataset", default="fake") + parser.add_argument("--items", nargs="*", default=[]) + parser.add_argument("--record", type=Path, default=None, help="Directory to write each request body to.") + parser.add_argument("--verbose", action="store_true") + return parser.parse_args(argv) + + +def _make_recorder(record_dir: Path): + record_dir.mkdir(parents=True, exist_ok=True) + numbers = itertools.count(1) + + def _record(record: dict) -> None: + safe_path = record["path"].replace("/", "_") + out = record_dir / f"{next(numbers):03d}-{record['method']}-{safe_path}.json" + out.write_text(json.dumps(record.get("json"), indent=2, default=str)) + + return _record + + +def main(argv: list[str]) -> int: + args = _parse_args(argv) + server = FakeLangfuse(port=args.port, verbose=args.verbose) + server.dataset_name = args.dataset + server.items = [_canned_item(item_id, args.dataset) for item_id in args.items] + + recorder = _make_recorder(args.record) if args.record else None + + def _on_request(record: dict) -> None: + print(json.dumps(record, indent=2, default=str)) + if recorder is not None: + recorder(record) + + server.on_request = _on_request + with server: + print(f"fake Langfuse server listening on {server.base_url}", file=sys.stderr) + try: + signal.pause() + except KeyboardInterrupt: + pass + return 0 + + +if __name__ == "__main__": + raise SystemExit(main(sys.argv[1:])) diff --git a/packages/gooddata-eval/tests/conftest.py b/packages/gooddata-eval/tests/conftest.py index 3b5b04734..560b8ebca 100644 --- a/packages/gooddata-eval/tests/conftest.py +++ b/packages/gooddata-eval/tests/conftest.py @@ -3,7 +3,20 @@ import pytest +from tests._fake_langfuse import FakeLangfuse + @pytest.fixture def fixtures_dir() -> Path: return Path(__file__).parent / "fixtures" + + +@pytest.fixture +def fake_langfuse(monkeypatch: pytest.MonkeyPatch): + """A running fake Langfuse server with the real client's env vars pointed at it.""" + with FakeLangfuse() as server: + monkeypatch.setenv("LANGFUSE_BASE_URL", server.base_url) + monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk-fake") + monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk-fake") + monkeypatch.delenv("LANGFUSE_HOST", raising=False) + yield server diff --git a/packages/gooddata-eval/tests/test_fake_langfuse.py b/packages/gooddata-eval/tests/test_fake_langfuse.py new file mode 100644 index 000000000..a785efeb4 --- /dev/null +++ b/packages/gooddata-eval/tests/test_fake_langfuse.py @@ -0,0 +1,140 @@ +# (C) 2026 GoodData Corporation +"""Tests for the in-process fake Langfuse server used by the wire-level e2e tests.""" + +from __future__ import annotations + +import httpx +import pytest + +from tests._fake_langfuse import FakeLangfuse + + +@pytest.fixture +def server(): + with FakeLangfuse() as srv: + yield srv + + +@pytest.fixture +def client(server: FakeLangfuse): + with httpx.Client(base_url=server.base_url) as c: + yield c + + +def test_dataset_items_serves_configured_dataset(client: httpx.Client, server: FakeLangfuse) -> None: + server.dataset_name = "GDAI-2179" + server.items = [{"id": "item-1", "datasetName": "GDAI-2179", "input": {"question": "q"}}] + + resp = client.get("/api/public/dataset-items", params={"datasetName": "GDAI-2179"}) + + assert resp.status_code == 200 + body = resp.json() + assert body["data"] == server.items + assert body["meta"]["totalItems"] == 1 + + +def test_dataset_items_404s_other_dataset_name(client: httpx.Client, server: FakeLangfuse) -> None: + server.dataset_name = "GDAI-2179" + + resp = client.get("/api/public/dataset-items", params={"datasetName": "other"}) + + assert resp.status_code == 404 + + +def test_dataset_item_lookup_returns_dataset_id(client: httpx.Client, server: FakeLangfuse) -> None: + server.dataset_id = "ds-fake" + + resp = client.get("/api/public/dataset-items/item-1") + + assert resp.status_code == 200 + assert resp.json()["datasetId"] == "ds-fake" + + +def test_dataset_item_lookup_404s_missing_id(client: httpx.Client, server: FakeLangfuse) -> None: + server.missing = {"item-404"} + + resp = client.get("/api/public/dataset-items/item-404") + + assert resp.status_code == 404 + + +def test_observations_synthesises_root_and_children_for_session(client: httpx.Client, server: FakeLangfuse) -> None: + resp = client.get( + "/api/public/v2/observations", + params={"fromStartTime": "2026-01-01T00:00:00Z", "toStartTime": "2026-01-02T00:00:00Z", "sessionId": "conv-1"}, + ) + + assert resp.status_code == 200 + rows = resp.json()["data"] + assert len(rows) == 3 + root = next(r for r in rows if r["parentObservationId"] is None) + assert root["metadata"]["conversation_id"] == "conv-1" + assert sum(r.get("totalCost") or 0.0 for r in rows) == pytest.approx(0.03) + + +def test_observations_two_page_cursor_splits_root_and_children(client: httpx.Client, server: FakeLangfuse) -> None: + server.observations_pages = 2 + params = {"fromStartTime": "2026-01-01T00:00:00Z", "toStartTime": "2026-01-02T00:00:00Z", "sessionId": "conv-2"} + + page1 = client.get("/api/public/v2/observations", params=params) + assert page1.status_code == 200 + body1 = page1.json() + assert len(body1["data"]) == 1 + assert body1["data"][0]["parentObservationId"] is None + cursor = body1["meta"]["cursor"] + assert cursor == "page-2" + + page2 = client.get("/api/public/v2/observations", params={**params, "cursor": cursor}) + assert page2.status_code == 200 + body2 = page2.json() + assert len(body2["data"]) == 2 + assert all(r["parentObservationId"] is not None for r in body2["data"]) + assert body2["meta"] == {} + + +def test_observations_first_call_empty_then_rows(client: httpx.Client, server: FakeLangfuse) -> None: + server.first_observations_call_empty = True + params = {"fromStartTime": "2026-01-01T00:00:00Z", "toStartTime": "2026-01-02T00:00:00Z", "sessionId": "conv-3"} + + first = client.get("/api/public/v2/observations", params=params) + assert first.status_code == 200 + assert first.json() == {"data": [], "meta": {}} + + second = client.get("/api/public/v2/observations", params=params) + assert second.status_code == 200 + assert len(second.json()["data"]) == 3 + + +def test_otlp_traces_records_method_path_query_and_json(client: httpx.Client, server: FakeLangfuse) -> None: + body = {"resourceSpans": [{"spanId": "abc"}]} + + resp = client.post("/api/public/otel/v1/traces", json=body) + + assert resp.status_code == 200 + calls = server.calls("POST", "/api/public/otel/v1/traces") + assert len(calls) == 1 + call = calls[0] + assert call["method"] == "POST" + assert call["path"] == "/api/public/otel/v1/traces" + assert call["json"] == body + + +def test_otlp_traces_honours_configured_status_and_body(client: httpx.Client, server: FakeLangfuse) -> None: + server.otlp_status = 400 + server.otlp_body = {"partialSuccess": {"rejectedSpans": 1}} + + resp = client.post("/api/public/otel/v1/traces", json={}) + + assert resp.status_code == 400 + assert resp.json() == {"partialSuccess": {"rejectedSpans": 1}} + + +def test_scores_429_once_then_200(client: httpx.Client, server: FakeLangfuse) -> None: + server.scores_429_once = True + + first = client.post("/api/public/scores", json={"traceId": "t1", "name": "n", "value": 1.0, "dataType": "NUMERIC"}) + assert first.status_code == 429 + assert first.headers["Retry-After"] == "0" + + second = client.post("/api/public/scores", json={"traceId": "t1", "name": "n", "value": 1.0, "dataType": "NUMERIC"}) + assert second.status_code == 200 From 18073f873ccb712c1216fceea8464becf776aef5 Mon Sep 17 00:00:00 2001 From: Jan Tychtl Date: Wed, 9 Sep 2026 21:57:45 +0200 Subject: [PATCH 3/3] fix(gooddata-eval): fall back to the default score retry delay for negative or NaN Retry-After time.sleep rejects a negative or NaN delay, so a skewed or hostile Retry-After header would turn a throttled score write into an exception instead of a retry. Such values now use the default delay; the existing cap still bounds long waits. jira: trivial risk: low --- .../src/gooddata_eval/core/langfuse/client.py | 11 +++++++++-- .../tests/test_langfuse_client.py | 18 ++++++++++++++++++ 2 files changed, 27 insertions(+), 2 deletions(-) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py index a52402381..138719b4c 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py @@ -32,11 +32,18 @@ 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.""" + """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. + """ try: - return min(float(resp.headers.get("Retry-After", "")), _MAX_RETRY_DELAY) + asked_for = float(resp.headers.get("Retry-After", "")) except ValueError: return _DEFAULT_RETRY_DELAY + if not asked_for >= 0: + return _DEFAULT_RETRY_DELAY + return min(asked_for, _MAX_RETRY_DELAY) class _TraceListResult: diff --git a/packages/gooddata-eval/tests/test_langfuse_client.py b/packages/gooddata-eval/tests/test_langfuse_client.py index 49a59daa9..f1bf5f228 100644 --- a/packages/gooddata-eval/tests/test_langfuse_client.py +++ b/packages/gooddata-eval/tests/test_langfuse_client.py @@ -139,6 +139,24 @@ def test_a_retry_delay_the_server_did_not_name_falls_back_to_a_short_one(make_cl assert slept == [0.5] +@pytest.mark.parametrize("header", ["-5", "nan"]) +def test_a_retry_after_that_cannot_be_slept_falls_back_to_the_default(header): + # time.sleep rejects a negative or NaN delay, so either would turn a retry into an exception. + 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_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. + response = httpx.Response(429, headers={"Retry-After": "Wed, 21 Oct 2026 07:28:00 GMT"}) + assert client_module._retry_delay(response) == 0.5 + + def test_a_score_gives_up_after_three_attempts(make_client, monkeypatch): monkeypatch.setattr(client_module.time, "sleep", lambda _seconds: None) attempts = 0