diff --git a/tests/e2e/conftest.py b/tests/e2e/conftest.py index 82f2604d492..4c4c8fe735b 100644 --- a/tests/e2e/conftest.py +++ b/tests/e2e/conftest.py @@ -1,9 +1,9 @@ """Shared fixtures for all live e2e suites under tests/e2e/. -Design rule: skip on environment, fail on behavior. Live tests (marked `e2e`) -skip when no proxy answers; once a request reaches the proxy, behavior is -asserted. Pure unit coverage of the harness itself carries no `e2e` marker and -runs regardless of whether a proxy is up. +Design rule: hard failures only. Live tests (marked `e2e`) fail when no proxy +answers or when credentials/env are missing; they never skip. Pure unit coverage +of the harness itself carries no `e2e` marker and runs regardless of whether a +proxy is up. Lifecycle: the `resources` fixture maps the init -> run -> teardown contract (lifecycle.E2ECase) onto pytest - setup is init(), the test body is run(), and @@ -40,7 +40,7 @@ def pytest_configure(config: pytest.Config) -> None: def _liveness_reason(label: str, base_url: str) -> str | None: - """None if `base_url` answers its liveness probe, else a skip reason.""" + """None if `base_url` answers its liveness probe, else a failure reason.""" try: resp = requests.get(f"{base_url}/health/liveliness", timeout=5) except requests.RequestException as exc: @@ -51,10 +51,10 @@ def _liveness_reason(label: str, base_url: str) -> str | None: @functools.lru_cache(maxsize=1) -def _proxy_skip_reason() -> str | None: - """Probe the proxy once per session. None if it answers, else a skip reason. In - a split deployment the management/admin control plane is a separate service, so - require it too (when it differs) - else its tests would fail rather than skip.""" +def _proxy_fail_reason() -> str | None: + """Probe the proxy once per session. None if it answers, else a failure reason. + In a split deployment the management/admin control plane is a separate service, + so require it too when it differs.""" reason = _liveness_reason("proxy", PROXY_BASE_URL) if reason is not None: return reason @@ -64,19 +64,19 @@ def _proxy_skip_reason() -> str | None: def pytest_runtest_setup(item: pytest.Item) -> None: - """Skip `e2e`-marked tests unless a proxy answers its liveness probe. Unmarked - tests (unit coverage of the harness) don't touch the proxy, so they run even - when none is up.""" + """Hard-fail `e2e`-marked tests unless a proxy answers its liveness probe. + Unmarked tests (unit coverage of the harness) don't touch the proxy, so they + run even when none is up. Never skip for a missing proxy.""" if item.get_closest_marker("e2e") is None: return - reason = _proxy_skip_reason() + reason = _proxy_fail_reason() if reason is not None: - pytest.skip(reason) + pytest.fail(reason) def pytest_runtest_call(item: pytest.Item) -> None: - """Mark that an e2e test body actually ran (not skipped at setup). Skipped - sessions never reach this hook, so the session-finish cleanup can use it as a + """Mark that an e2e test body actually ran (setup passed). Sessions that fail + setup never reach this hook, so the session-finish cleanup can use it as a guard before truncating the spend-log DB. Tests under `tests/e2e/` without the `e2e` marker (pure unit coverage for the harness itself) never hit the proxy, so they must not arm the destructive DB truncate.""" @@ -87,12 +87,12 @@ def pytest_runtest_call(item: pytest.Item) -> None: def pytest_sessionfinish(session: pytest.Session, exitstatus: int) -> None: """Once the whole e2e session is done (all suites), truncate the spend logs so - the DB doesn't accumulate test rows. Skipped sessions (no live proxy, no test - actually executed) leave the DB alone so a `DATABASE_URL` pointing at a shared - instance is never wiped without an e2e run. Best-effort: a cleanup failure (no - DB reachable) must not fail the run. The spend_tracking dir goes on sys.path - only for this import and is removed after, so a broader `pytest tests/` run is - not left with a mutated path.""" + the DB doesn't accumulate test rows. Sessions where no e2e test body ran leave + the DB alone so a `DATABASE_URL` pointing at a shared instance is never wiped + without an e2e run. Best-effort: a cleanup failure (no DB reachable) must not + fail the run. The spend_tracking dir goes on sys.path only for this import and + is removed after, so a broader `pytest tests/` run is not left with a mutated + path.""" if not session.stash.get(_E2E_TEST_RAN, False): return spend_dir = str(Path(__file__).parent / "spend_tracking") @@ -102,7 +102,7 @@ def pytest_sessionfinish(session: pytest.Session, exitstatus: int) -> None: reset_spend_logs() except Exception as exc: # noqa: BLE001 - cleanup is best-effort - print(f"spend-log cleanup skipped: {exc}") + print(f"spend-log cleanup best-effort failed: {exc}") finally: if spend_dir in sys.path: sys.path.remove(spend_dir) @@ -112,7 +112,7 @@ def pytest_sessionfinish(session: pytest.Session, exitstatus: int) -> None: remediate(session) except Exception as exc: # noqa: BLE001 - remediation is best-effort - print(f"devin remediation skipped: {exc}") + print(f"devin remediation best-effort failed: {exc}") @pytest.fixture diff --git a/tests/e2e/e2e_http.py b/tests/e2e/e2e_http.py index 7b8d3045d3b..32005faed4b 100644 --- a/tests/e2e/e2e_http.py +++ b/tests/e2e/e2e_http.py @@ -107,13 +107,15 @@ class ProbeResult(BaseModel): class StreamingResponse(BaseModel): """Raw outcome for calls whose body is provider-native or streamed: status, the - x-litellm-call-id header (== SpendLogs.request_id), the content-type (which - tells streaming `text/event-stream` from non-streaming `application/json`), and - the body. Used by passthrough and streaming, where one validated JSON model - does not fit.""" + x-litellm-call-id header, the x-litellm-response-cost header (StandardLogging + response_cost), the content-type (which tells streaming `text/event-stream` from + non-streaming `application/json`), and the body. SpendLogs.request_id is the + completion body id, not call_id. Used by passthrough and streaming, where one + validated JSON model does not fit.""" status_code: int call_id: str | None = None # x-litellm-call-id header + response_cost: float | None = None # x-litellm-response-cost header content_type: str | None = None body: str chunks: int = 0 # streamed events (0 for non-streaming) @@ -260,13 +262,25 @@ def probe( return ProbeResult(status_code=resp.status_code, body=resp.text) +def _parse_response_cost(resp: requests.Response) -> float | None: + raw = _hdr(resp, "x-litellm-response-cost") + if raw is None or raw == "": + return None + try: + return float(raw) + except ValueError: + return None + + def _streaming_outcome(resp: requests.Response, stream: bool) -> StreamingResponse: call_id = _hdr(resp, "x-litellm-call-id") + response_cost = _parse_response_cost(resp) content_type = _hdr(resp, "content-type") if not stream or not (200 <= resp.status_code < 300): return StreamingResponse( status_code=resp.status_code, call_id=call_id, + response_cost=response_cost, content_type=content_type, body=resp.text, ) @@ -275,6 +289,7 @@ def _streaming_outcome(resp: requests.Response, stream: bool) -> StreamingRespon return StreamingResponse( status_code=resp.status_code, call_id=call_id, + response_cost=response_cost, content_type=content_type, body="", chunks=chunks, diff --git a/tests/e2e/logging/conftest.py b/tests/e2e/logging/conftest.py index 40c19aefca7..567355f23ac 100644 --- a/tests/e2e/logging/conftest.py +++ b/tests/e2e/logging/conftest.py @@ -1,37 +1,43 @@ -"""Fixtures for the Datadog logging suite. +"""Fixtures for the logging e2e suite. -These tests drive the Datadog batch-send path (#25663) directly against the real -Datadog logs intake with synthetic events - no LLM calls, no proxy, no log -read-back - so they need only the shipping credentials DD_API_KEY + DD_SITE -(DD_SERVICE is an optional tag). No Datadog Application key is required, and they -skip when the shipping credentials are absent from the environment. +Missing proxy, provider keys, or integration credentials are hard failures. +Never pytest.skip from this suite for environment gaps. """ +from __future__ import annotations + import os import pytest -from logging_client import LoggingClient, build_logging_client +from logging_client import LangfuseCreds, LoggingClient, build_logging_client, load_langfuse_creds def pytest_configure(config: pytest.Config) -> None: config.addinivalue_line( "markers", - "covers: registry cell a test covers, e.g. logging.datadog.success.writes_object", + "covers: registry cell a test covers, e.g. logging.langfuse.success.logs_spend", ) @pytest.fixture(scope="session") def client() -> LoggingClient: """The logging suite's client: holds the shared Gateway so `resources` / - `scoped_key` clean up keys, and adds `/metrics` scraping.""" + `scoped_key` clean up keys and teams, and adds `/metrics` scraping plus + Langfuse read-back.""" return build_logging_client() @pytest.fixture def datadog_creds() -> None: - """Gate the suite on the Datadog shipping credentials. The DataDogLogger is built - inside each async test, not here, because its __init__ schedules a periodic-flush - task via asyncio.create_task and so needs a running event loop.""" + """Require Datadog shipping credentials. Hard-fail when absent; never skip.""" if not (os.getenv("DD_API_KEY") and os.getenv("DD_SITE")): - pytest.skip("set DD_API_KEY and DD_SITE to run the Datadog logging suite") + pytest.fail( + "Datadog e2e requires DD_API_KEY and DD_SITE; missing credentials is a hard failure, not a skip" + ) + + +@pytest.fixture(scope="session") +def langfuse_creds() -> LangfuseCreds: + """Require real Langfuse cloud credentials for team callback + trace poll.""" + return load_langfuse_creds() diff --git a/tests/e2e/logging/logging_client.py b/tests/e2e/logging/logging_client.py index a3213fbdb00..06b219fc6e2 100644 --- a/tests/e2e/logging/logging_client.py +++ b/tests/e2e/logging/logging_client.py @@ -1,48 +1,571 @@ -"""Client for the logging e2e suite: drive traffic and scrape the proxy's -Prometheus ``/metrics`` endpoint. +"""Client for the logging e2e suite: team/key/org-scoped Langfuse OTEL callbacks, +chat (including tools), Prometheus scrape, and Langfuse observation read-back. -Holds the shared Gateway so the ``resources`` fixture cleans up keys it creates. -``/metrics`` is exposed as plaintext (not a typed JSON body), so scraping goes -through ``transport.probe`` and returns the raw exposition text for a Prometheus -parser to read. +Holds the shared Gateway so the ``resources`` fixture cleans up keys, teams, +users, orgs, and models it creates. External Langfuse reads go through +``e2e_http`` (the only module allowed to call ``requests.*``). + +Uses the ``langfuse_otel`` callback (OTLP to ``{host}/api/public/otel``), not +the classic ``langfuse`` SDK callback. OTEL generations land as name +``litellm_request``; correlate by unique prompt marker and ``user_api_key_alias`` +in metadata. Spend is on ``calculatedTotalCost`` (StandardLogging response_cost). """ from __future__ import annotations +import base64 +import json +import os +import time from dataclasses import dataclass +from typing import Literal +import pytest +from pydantic import BaseModel, ConfigDict, Field + +from e2e_config import POLL_INTERVAL, POLL_TIMEOUT from e2e_gateway import Gateway, build_gateway -from e2e_http import NoBody, unwrap -from models import ChatBody, ChatMessage, ChatResponse, KeyGenerateBody +from e2e_http import ( + URL, + AuthHeaders, + NoBody, + StreamingResponse, + Success, + get, + unwrap, +) +from models import ( + ChatBody, + ChatMessage, + ChatResponse, + ChatTool, + ChatToolFunction, + KeyGenerateBody, + KeyLoggingCallback, + KeyLoggingCallbackVars, + KeyMetadata, + LiteLLMParamsBody, + OrgDeleteBody, + OrgNewBody, + OrgNewResponse, + SpendLogRow, + TeamDeleteBody, + TeamNewBody, + TeamNewResponse, + UserDeleteBody, + UserNewBody, + UserNewResponse, +) + +# Deliberately invalid *upstream provider* key for failure-path tests. +# Not a LiteLLM virtual key; OpenAI must reject it after the proxy accepts the call. +INVALID_UPSTREAM_API_KEY = "sk-upstream-invalid-for-langfuse-e2e-only" + +WEATHER_TOOL = ChatTool( + type="function", + function=ChatToolFunction( + name="get_weather", + description="Get the current weather for a city", + parameters={ + "type": "object", + "properties": {"city": {"type": "string"}}, + "required": ["city"], + }, + ), +) + + +class TeamCallbackBody(BaseModel): + callback_name: Literal["langfuse_otel", "langfuse", "langsmith", "gcs"] + callback_type: Literal["success", "failure", "success_and_failure"] + callback_vars: dict[str, str] + + +class TeamCallbackResponse(BaseModel): + model_config = ConfigDict(extra="ignore") + + status: str + + +class GuardrailLitellmParams(BaseModel): + guardrail: str + mode: str + default_on: bool = False + rules: list[dict[str, object]] | None = None + default_action: str | None = None + on_disallowed_action: str | None = None + + +class GuardrailSpec(BaseModel): + guardrail_name: str + litellm_params: GuardrailLitellmParams + + +class CreateGuardrailBody(BaseModel): + guardrail: GuardrailSpec + + +class CreateGuardrailResponse(BaseModel): + model_config = ConfigDict(extra="ignore") + + guardrail_id: str | None = None + guardrail_name: str | None = None + + +class LangfuseObservation(BaseModel): + model_config = ConfigDict(extra="ignore", populate_by_name=True) + + id: str + trace_id: str | None = Field(default=None, alias="traceId") + name: str | None = None + type: str | None = None + calculated_total_cost: float | None = Field(default=None, alias="calculatedTotalCost") + level: str | None = None + input: object | None = None + output: object | None = None + metadata: object | None = None + usage: object | None = None + usage_details: object | None = Field(default=None, alias="usageDetails") + model: str | None = None + + +class LangfuseObservationList(BaseModel): + model_config = ConfigDict(extra="ignore") + + data: list[LangfuseObservation] = [] + + +class LangfuseListParams(BaseModel): + model_config = ConfigDict(populate_by_name=True) + + limit: int = 100 + trace_id: str | None = Field(default=None, alias="traceId") + name: str | None = None + from_start_time: str | None = Field(default=None, alias="fromStartTime") + + +@dataclass(frozen=True, slots=True) +class LangfuseCreds: + public_key: str + secret_key: str + host: str + + @property + def auth_headers(self) -> AuthHeaders: + token = base64.b64encode(f"{self.public_key}:{self.secret_key}".encode()).decode() + return AuthHeaders(authorization=f"Basic {token}") + + def callback_vars(self) -> dict[str, str]: + return { + "langfuse_public_key": self.public_key, + "langfuse_secret_key": self.secret_key, + "langfuse_host": self.host, + } + + def key_logging_metadata(self) -> KeyMetadata: + return KeyMetadata( + logging=[ + KeyLoggingCallback( + callback_name="langfuse_otel", + callback_type="success_and_failure", + callback_vars=KeyLoggingCallbackVars( + langfuse_public_key=self.public_key, + langfuse_secret_key=self.secret_key, + langfuse_host=self.host, + ), + ) + ] + ) + + +def load_langfuse_creds() -> LangfuseCreds: + public_key = os.getenv("LANGFUSE_PUBLIC_KEY") + secret_key = os.getenv("LANGFUSE_SECRET_KEY") + host = (os.getenv("LANGFUSE_BASE_URL") or os.getenv("LANGFUSE_HOST") or "").rstrip("/") + if not (public_key and secret_key and host): + pytest.fail( + "Langfuse e2e requires LANGFUSE_PUBLIC_KEY, LANGFUSE_SECRET_KEY, and " + "LANGFUSE_BASE_URL (or LANGFUSE_HOST); missing credentials is a hard failure, not a skip" + ) + return LangfuseCreds(public_key=public_key, secret_key=secret_key, host=host) + + +def observation_spend(obs: LangfuseObservation) -> float | None: + """Langfuse calculatedTotalCost is populated from StandardLogging response_cost.""" + return obs.calculated_total_cost + + +def costs_agree(expected: float, actual: float, *, rel_tol: float = 0.05) -> bool: + """Costs agree within 5% relative (or 1e-9 absolute for near-zero).""" + return abs(expected - actual) <= max(1e-9, abs(expected) * rel_tol) + + +def completion_response_id(body: str) -> str | None: + """SpendLogs.request_id is the chat completion body id, not x-litellm-call-id.""" + if not body or body == "": + return None + try: + parsed = json.loads(body) + except json.JSONDecodeError: + return None + if not isinstance(parsed, dict): + return None + raw = parsed.get("id") + return raw if isinstance(raw, str) and raw else None + + +def _matches_run(obs: LangfuseObservation, *, key_alias: str, prompt_marker: str) -> bool: + """Match a Langfuse generation for this run. + + langfuse_otel names generations ``litellm_request`` (not ``litellm:{alias}``). + Prefer the unique prompt marker in input; fall back to key alias in metadata + (user_api_key_alias) or the classic SDK generation name. + """ + if prompt_marker and prompt_marker in json.dumps(obs.input, default=str): + return True + meta_blob = json.dumps(obs.metadata, default=str) if obs.metadata is not None else "" + if key_alias and key_alias in meta_blob: + return True + if obs.name == f"litellm:{key_alias}": + return True + return False + + +def observation_mentions_tool(obs: LangfuseObservation, tool_name: str) -> bool: + blob = json.dumps( + {"input": obs.input, "output": obs.output, "metadata": obs.metadata}, + default=str, + ) + return tool_name in blob + + +def observation_has_guardrail(obs: LangfuseObservation, *, guardrail_name: str) -> bool: + blob = json.dumps(obs.metadata, default=str) if obs.metadata is not None else "" + if guardrail_name in blob or "guardrail" in blob.lower(): + return True + if obs.name is not None and "guardrail" in obs.name.lower(): + return True + return False @dataclass(frozen=True, slots=True) class LoggingClient: gateway: Gateway - def key_with_alias(self, alias: str, *, models: list[str]) -> str: + def key_with_alias( + self, + alias: str, + *, + models: list[str], + team_id: str | None = None, + user_id: str | None = None, + organization_id: str | None = None, + metadata: KeyMetadata | None = None, + ) -> str: return self.gateway.generate_key( - KeyGenerateBody(key_alias=alias, models=models, user_id=f"e2e-{alias}") + KeyGenerateBody( + key_alias=alias, + models=models, + user_id=user_id or f"e2e-{alias}", + team_id=team_id, + organization_id=organization_id, + metadata=metadata, + ) ) def delete_key(self, key: str) -> None: self.gateway.delete_key(key) + def create_team( + self, + alias: str, + *, + models: list[str], + organization_id: str | None = None, + ) -> str: + return unwrap( + self.gateway.transport.post( + "/team/new", + headers=self.gateway.transport.master, + json=TeamNewBody( + team_alias=alias, + models=models, + organization_id=organization_id, + ), + response_type=TeamNewResponse, + ) + ).team_id + + def delete_team(self, team_id: str) -> None: + _ = self.gateway.transport.post( + "/team/delete", + headers=self.gateway.transport.master, + json=TeamDeleteBody(team_ids=[team_id]), + response_type=NoBody, + ) + + def create_user(self, *, user_email: str, user_id: str | None = None) -> str: + return unwrap( + self.gateway.transport.post( + "/user/new", + headers=self.gateway.transport.master, + json=UserNewBody( + user_email=user_email, + user_role="internal_user", + user_id=user_id, + ), + response_type=UserNewResponse, + ) + ).user_id + + def delete_user(self, user_id: str) -> None: + _ = self.gateway.transport.post( + "/user/delete", + headers=self.gateway.transport.master, + json=UserDeleteBody(user_ids=[user_id]), + response_type=NoBody, + ) + + def create_org(self, alias: str, *, models: list[str]) -> str: + return unwrap( + self.gateway.transport.post( + "/organization/new", + headers=self.gateway.transport.master, + json=OrgNewBody(organization_alias=alias, models=models), + response_type=OrgNewResponse, + ) + ).organization_id + + def delete_org(self, organization_id: str) -> None: + _ = self.gateway.transport.delete( + "/organization/delete", + headers=self.gateway.transport.master, + json=OrgDeleteBody(organization_ids=[organization_id]), + response_type=NoBody, + ) + + def add_team_langfuse_callback( + self, + team_id: str, + creds: LangfuseCreds, + *, + callback_type: Literal["success", "failure", "success_and_failure"] = "success_and_failure", + ) -> None: + response = unwrap( + self.gateway.transport.post( + f"/team/{team_id}/callback", + headers=self.gateway.transport.master, + json=TeamCallbackBody( + callback_name="langfuse_otel", + callback_type=callback_type, + callback_vars=creds.callback_vars(), + ), + response_type=TeamCallbackResponse, + ) + ) + assert response.status == "success", ( + f"POST /team/{team_id}/callback must return status=success; got {response.status!r}" + ) + + def create_tool_permission_guardrail(self, name: str, *, allowed_tool: str) -> str: + """Register a tool_permission guardrail that allows one tool and denies the rest.""" + response = unwrap( + self.gateway.transport.post( + "/guardrails", + headers=self.gateway.transport.master, + json=CreateGuardrailBody( + guardrail=GuardrailSpec( + guardrail_name=name, + litellm_params=GuardrailLitellmParams( + guardrail="tool_permission", + mode="post_call", + default_on=False, + default_action="deny", + on_disallowed_action="block", + rules=[ + { + "id": "allow-named-tool", + "tool_name": allowed_tool, + "decision": "allow", + } + ], + ), + ) + ), + response_type=CreateGuardrailResponse, + ) + ) + guardrail_id = response.guardrail_id + assert guardrail_id, f"create guardrail returned no id: {response!r}" + return guardrail_id + + def delete_guardrail(self, guardrail_id: str) -> None: + _ = self.gateway.transport.delete( + f"/guardrails/{guardrail_id}", + headers=self.gateway.transport.master, + json=NoBody(), + response_type=NoBody, + ) + + def create_model(self, model_name: str, litellm_params: LiteLLMParamsBody) -> str: + return self.gateway.create_model(model_name, litellm_params) + + def delete_model(self, model_id: str) -> None: + self.gateway.delete_model(model_id) + def chat(self, key: str, model: str, text: str) -> ChatResponse: return unwrap( self.gateway.chat( key, ChatBody( model=model, - messages=[ChatMessage(role="user", content=text)], - max_tokens=64, + messages=[ChatMessage(role="user", content=text)], + max_tokens=64, ), ) ) + def chat_raw( + self, + key: str, + model: str, + text: str, + *, + stream: bool = False, + tools: list[ChatTool] | None = None, + tool_choice: str | None = None, + guardrails: list[str] | None = None, + max_tokens: int = 64, + ) -> StreamingResponse: + body = ChatBody( + model=model, + messages=[ChatMessage(role="user", content=text)], + max_tokens=max_tokens, + stream=stream, + tools=tools, + tool_choice=tool_choice, + guardrails=guardrails, + ) + if stream: + return self.gateway.chat_stream(key, body) + return self.gateway.transport.send( + "/chat/completions", + headers=self.gateway.transport.bearer(key), + json=body, + ) + def scrape_metrics(self) -> str: return self.gateway.probe("/metrics", params=NoBody()).body + def poll_proxy_spend_for_key( + self, + key: str, + *, + response_id: str | None = None, + require_positive_spend: bool = True, + ) -> SpendLogRow | None: + """Poll /spend/logs by virtual key. + + When ``response_id`` is set, only that SpendLogs.request_id may match. + When unset, any positive-spend row for the key is accepted. Never falls + back to an unmatched row; missing match returns None. + """ + + def _matches(row: SpendLogRow) -> bool: + if response_id is not None and row.request_id != response_id: + return False + if require_positive_spend and not (row.spend is not None and row.spend > 0): + return False + return True + + rows = self.gateway.poll_logs_for_key( + key, min_rows=1, predicate=lambda rs: any(_matches(r) for r in rs) + ) + for row in rows: + if _matches(row): + return row + return None + + def list_langfuse_observations( + self, + creds: LangfuseCreds, + *, + trace_id: str | None = None, + name: str | None = None, + from_start_time: str | None = None, + ) -> list[LangfuseObservation]: + result = get( + URL(f"{creds.host}/api/public/observations"), + headers=creds.auth_headers, + params=LangfuseListParams( + limit=100, + trace_id=trace_id, + name=name, + from_start_time=from_start_time, + ), + response_type=LangfuseObservationList, + timeout=30.0, + ) + match result: + case Success(data=page): + return page.data + case _: + return [] + + def find_langfuse_observation( + self, + creds: LangfuseCreds, + *, + key_alias: str, + prompt_marker: str, + ) -> LangfuseObservation | None: + # langfuse_otel generations are named litellm_request; classic SDK used + # litellm:{key_alias}. Search both, then a recent unfiltered page. + for name in ("litellm_request", f"litellm:{key_alias}"): + for obs in self.list_langfuse_observations(creds, name=name): + if _matches_run(obs, key_alias=key_alias, prompt_marker=prompt_marker): + return obs + for obs in self.list_langfuse_observations(creds): + if _matches_run(obs, key_alias=key_alias, prompt_marker=prompt_marker): + return obs + return None + + def poll_langfuse_observation( + self, + creds: LangfuseCreds, + *, + key_alias: str, + prompt_marker: str, + require_positive_cost: bool = False, + ) -> LangfuseObservation | None: + deadline = time.monotonic() + POLL_TIMEOUT + last: LangfuseObservation | None = None + while time.monotonic() < deadline: + last = self.find_langfuse_observation( + creds, key_alias=key_alias, prompt_marker=prompt_marker + ) + if last is not None: + cost = observation_spend(last) + if not require_positive_cost or (cost is not None and cost > 0): + return last + time.sleep(POLL_INTERVAL) + return last + + def poll_langfuse_trace_observations( + self, + creds: LangfuseCreds, + *, + key_alias: str, + prompt_marker: str, + ) -> list[LangfuseObservation]: + """Generation plus any sibling/child observations (guardrail spans, etc.).""" + gen = self.poll_langfuse_observation( + creds, key_alias=key_alias, prompt_marker=prompt_marker + ) + if gen is None or not gen.trace_id: + return [] if gen is None else [gen] + return self.list_langfuse_observations(creds, trace_id=gen.trace_id) or [gen] + def build_logging_client() -> LoggingClient: return LoggingClient(gateway=build_gateway()) diff --git a/tests/e2e/logging/test_langfuse_e2e.py b/tests/e2e/logging/test_langfuse_e2e.py new file mode 100644 index 00000000000..d014b5d8291 --- /dev/null +++ b/tests/e2e/logging/test_langfuse_e2e.py @@ -0,0 +1,534 @@ +"""Live e2e: Langfuse OTEL logs_spend for registry cells in logging.yaml P0. + +Registry cells: +- logging.langfuse.success.logs_spend (exercised_on chat_completions, messages, embeddings) +- logging.langfuse.failure.logs_spend (exercised_on chat_completions, messages) +- logging.langfuse.stream.logs_spend (exercised_on chat_completions, messages) + +Integration under test is ``langfuse_otel`` (OTLP to Langfuse), not the classic +``langfuse`` SDK callback. StandardLoggingPayload.response_cost is the spend +source of truth. Generations are named ``litellm_request``; correlate by unique +prompt marker and user_api_key_alias in metadata. + +Dynamic credentials by product surface: +- team: POST /team/{id}/callback with callback_name=langfuse_otel +- user/key: key metadata.logging with callback_name=langfuse_otel +- org: organization + team under it + team callback (no org-level callback API) + +Extra success paths assert tool calls and applied guardrails land on the trace. +""" + +from __future__ import annotations + +import json + +import pytest + +from e2e_config import unique_marker +from e2e_http import StreamingResponse, require_successful_call +from lifecycle import ResourceManager +from logging_client import ( + INVALID_UPSTREAM_API_KEY, + WEATHER_TOOL, + LangfuseCreds, + LoggingClient, + completion_response_id, + costs_agree, + observation_has_guardrail, + observation_mentions_tool, + observation_spend, +) +from models import LiteLLMParamsBody + +pytestmark = pytest.mark.e2e + +DRIVER_MODEL = "gemini-2.5-flash" +FAIL_BACKEND = "openai/gpt-4o-mini" + + +def _json_blob(value: object) -> str: + return json.dumps(value, default=str) + + +def _assert_logs_spend( + client: LoggingClient, + *, + key: str, + outcome: StreamingResponse, + obs_cost: float | None, + scope: str, + require_positive: bool = True, +) -> None: + """logs_spend: Langfuse cost matches StandardLogging response_cost and proxy spend. + + Non-stream responses expose response_cost on x-litellm-response-cost. Streaming + sends headers before final cost is known, so stream paths rely on /spend/logs. + """ + if not require_positive: + assert obs_cost is not None, ( + f"{scope}: failure path must still track spend (0 is fine); cost={obs_cost!r}" + ) + return + + assert obs_cost is not None and obs_cost > 0, ( + f"{scope}: Langfuse must log positive spend; calculatedTotalCost={obs_cost!r}" + ) + # Stream responses send headers before final cost is known, so the cost header + # is often absent; non-stream must always expose x-litellm-response-cost. + if not outcome.is_streaming: + assert outcome.response_cost is not None and outcome.response_cost > 0, ( + f"{scope}: proxy must return positive x-litellm-response-cost; " + f"got {outcome.response_cost!r}" + ) + assert costs_agree(outcome.response_cost, obs_cost), ( + f"{scope}: Langfuse cost {obs_cost!r} disagrees with " + f"x-litellm-response-cost {outcome.response_cost!r}" + ) + elif outcome.response_cost is not None and outcome.response_cost > 0: + assert costs_agree(outcome.response_cost, obs_cost), ( + f"{scope}: Langfuse cost {obs_cost!r} disagrees with " + f"x-litellm-response-cost {outcome.response_cost!r}" + ) + spend_row = client.poll_proxy_spend_for_key( + key, + response_id=completion_response_id(outcome.body), + require_positive_spend=True, + ) + assert spend_row is not None and spend_row.spend is not None and spend_row.spend > 0, ( + f"{scope}: proxy /spend/logs never produced a positive spend row for key" + ) + assert costs_agree(spend_row.spend, obs_cost), ( + f"{scope}: Langfuse cost {obs_cost!r} disagrees with proxy spend " + f"{spend_row.spend!r} (request_id={spend_row.request_id!r})" + ) + + +class TestLangfuseTeamLogging: + """Team-scoped callback via POST /team/{id}/callback.""" + + def _team_key( + self, + client: LoggingClient, + resources: ResourceManager, + creds: LangfuseCreds, + *, + models: list[str], + organization_id: str | None = None, + ) -> tuple[str, str, str]: + marker = unique_marker() + key_alias = f"e2e-lf-team-key-{marker}" + team_id = client.create_team( + f"e2e-lf-team-{marker}", + models=models, + organization_id=organization_id, + ) + resources.defer(lambda: client.delete_team(team_id)) + client.add_team_langfuse_callback(team_id, creds) + key = client.key_with_alias(key_alias, models=models, team_id=team_id) + resources.defer(lambda: client.delete_key(key)) + return team_id, key, key_alias + + @pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["chat_completions"]) + def test_success_logs_spend( + self, + client: LoggingClient, + resources: ResourceManager, + langfuse_creds: LangfuseCreds, + ) -> None: + _, key, key_alias = self._team_key( + client, resources, langfuse_creds, models=[DRIVER_MODEL] + ) + prompt_marker = unique_marker() + outcome = client.chat_raw( + key, DRIVER_MODEL, f"reply with one word only {prompt_marker}" + ) + require_successful_call(outcome) + + obs = client.poll_langfuse_observation( + langfuse_creds, + key_alias=key_alias, + prompt_marker=prompt_marker, + require_positive_cost=True, + ) + assert obs is not None, ( + f"team scope: Langfuse never received generation for key_alias={key_alias!r}" + ) + _assert_logs_spend( + client, + key=key, + outcome=outcome, + obs_cost=observation_spend(obs), + scope="team-success", + ) + + @pytest.mark.covers("logging.langfuse.failure.logs_spend", exercised_on=["chat_completions"]) + def test_failure_logs_spend( + self, + client: LoggingClient, + resources: ResourceManager, + langfuse_creds: LangfuseCreds, + ) -> None: + """Provider-auth failure still ships a Langfuse observation with spend tracked. + + Uses a throwaway deployment whose upstream OpenAI key is + INVALID_UPSTREAM_API_KEY (not a LiteLLM virtual key). + """ + prompt_marker = unique_marker() + model_name = f"e2e-lf-fail-{prompt_marker}" + model_id = client.create_model( + model_name, + LiteLLMParamsBody(model=FAIL_BACKEND, api_key=INVALID_UPSTREAM_API_KEY), + ) + resources.defer(lambda: client.delete_model(model_id)) + + _, key, key_alias = self._team_key( + client, resources, langfuse_creds, models=[model_name] + ) + outcome = client.chat_raw(key, model_name, f"this must fail {prompt_marker}") + assert not outcome.ok, ( + f"expected upstream provider failure for {INVALID_UPSTREAM_API_KEY!r}, " + f"got {outcome.status_code}: {outcome.body[:200]}" + ) + + obs = client.poll_langfuse_observation( + langfuse_creds, + key_alias=key_alias, + prompt_marker=prompt_marker, + require_positive_cost=False, + ) + assert obs is not None, ( + f"team failure path: Langfuse never received generation for key_alias={key_alias!r}" + ) + _assert_logs_spend( + client, + key=key, + outcome=outcome, + obs_cost=observation_spend(obs), + scope="team-failure", + require_positive=False, + ) + + @pytest.mark.covers("logging.langfuse.stream.logs_spend", exercised_on=["chat_completions"]) + def test_stream_logs_spend( + self, + client: LoggingClient, + resources: ResourceManager, + langfuse_creds: LangfuseCreds, + ) -> None: + _, key, key_alias = self._team_key( + client, resources, langfuse_creds, models=[DRIVER_MODEL] + ) + prompt_marker = unique_marker() + outcome = client.chat_raw( + key, DRIVER_MODEL, f"reply with one word only {prompt_marker}", stream=True + ) + require_successful_call(outcome) + assert outcome.is_streaming + assert outcome.chunks > 0 + + obs = client.poll_langfuse_observation( + langfuse_creds, + key_alias=key_alias, + prompt_marker=prompt_marker, + require_positive_cost=True, + ) + assert obs is not None + # Streamed body is elided; correlate cost via header + key spend row. + _assert_logs_spend( + client, + key=key, + outcome=outcome, + obs_cost=observation_spend(obs), + scope="team-stream", + ) + + @pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["chat_completions"]) + def test_tool_calls_logged_with_cost( + self, + client: LoggingClient, + resources: ResourceManager, + langfuse_creds: LangfuseCreds, + ) -> None: + _, key, key_alias = self._team_key( + client, resources, langfuse_creds, models=[DRIVER_MODEL] + ) + prompt_marker = unique_marker() + outcome = client.chat_raw( + key, + DRIVER_MODEL, + f"Use get_weather for Paris. marker={prompt_marker}", + tools=[WEATHER_TOOL], + tool_choice="required", + max_tokens=128, + ) + require_successful_call(outcome) + assert "get_weather" in outcome.body or "tool_calls" in outcome.body, ( + f"gateway response must include a tool call; body={outcome.body[:300]}" + ) + + obs = client.poll_langfuse_observation( + langfuse_creds, + key_alias=key_alias, + prompt_marker=prompt_marker, + require_positive_cost=True, + ) + assert obs is not None + assert observation_mentions_tool(obs, "get_weather"), ( + f"Langfuse generation must record the tool; name={obs.name!r} " + f"input={str(obs.input)[:200]} output={str(obs.output)[:200]}" + ) + _assert_logs_spend( + client, + key=key, + outcome=outcome, + obs_cost=observation_spend(obs), + scope="team-tools", + ) + + @pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["chat_completions"]) + def test_tool_permission_guardrail_logged( + self, + client: LoggingClient, + resources: ResourceManager, + langfuse_creds: LangfuseCreds, + ) -> None: + """tool_permission post_call guardrail must appear on the Langfuse trace + (StandardLogging guardrail_information -> Langfuse guardrail span).""" + marker = unique_marker() + guardrail_name = f"e2e-lf-tool-perm-{marker}" + guardrail_id = client.create_tool_permission_guardrail( + guardrail_name, allowed_tool="get_weather" + ) + resources.defer(lambda: client.delete_guardrail(guardrail_id)) + + _, key, key_alias = self._team_key( + client, resources, langfuse_creds, models=[DRIVER_MODEL] + ) + prompt_marker = unique_marker() + outcome = client.chat_raw( + key, + DRIVER_MODEL, + f"Use get_weather for Berlin. marker={prompt_marker}", + tools=[WEATHER_TOOL], + tool_choice="required", + guardrails=[guardrail_name], + max_tokens=128, + ) + require_successful_call(outcome) + + observations = client.poll_langfuse_trace_observations( + langfuse_creds, key_alias=key_alias, prompt_marker=prompt_marker + ) + assert observations, ( + f"team+guardrail: no Langfuse observations for key_alias={key_alias!r}" + ) + gen = next( + ( + o + for o in observations + if prompt_marker in _json_blob(o.input) + or key_alias in _json_blob(o.metadata) + or o.name in (f"litellm:{key_alias}", "litellm_request") + ), + observations[0], + ) + _assert_logs_spend( + client, + key=key, + outcome=outcome, + obs_cost=observation_spend(gen), + scope="team-guardrail", + ) + assert any( + observation_has_guardrail(o, guardrail_name=guardrail_name) + or (o.name is not None and "guardrail" in o.name.lower()) + for o in observations + ), ( + f"Langfuse trace must include applied guardrail {guardrail_name!r}; " + f"observation names={[o.name for o in observations]}" + ) + + +class TestLangfuseUserKeyLogging: + """User-owned key with metadata.logging (key-level dynamic Langfuse credentials). + + Product surface: key metadata.logging on /key/generate, not a separate + /user/.../callback route. The key is bound to a real /user/new user_id. + """ + + def _user_key( + self, + client: LoggingClient, + resources: ResourceManager, + creds: LangfuseCreds, + *, + models: list[str], + ) -> tuple[str, str, str]: + marker = unique_marker() + key_alias = f"e2e-lf-user-key-{marker}" + user_id = client.create_user( + user_email=f"e2e-lf-user-{marker}@example.com", + user_id=f"e2e-lf-user-{marker}", + ) + resources.defer(lambda: client.delete_user(user_id)) + key = client.key_with_alias( + key_alias, + models=models, + user_id=user_id, + metadata=creds.key_logging_metadata(), + ) + resources.defer(lambda: client.delete_key(key)) + return user_id, key, key_alias + + @pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["chat_completions"]) + def test_success_logs_spend( + self, + client: LoggingClient, + resources: ResourceManager, + langfuse_creds: LangfuseCreds, + ) -> None: + user_id, key, key_alias = self._user_key( + client, resources, langfuse_creds, models=[DRIVER_MODEL] + ) + prompt_marker = unique_marker() + outcome = client.chat_raw( + key, DRIVER_MODEL, f"reply with one word only {prompt_marker}" + ) + require_successful_call(outcome) + + obs = client.poll_langfuse_observation( + langfuse_creds, + key_alias=key_alias, + prompt_marker=prompt_marker, + require_positive_cost=True, + ) + assert obs is not None, ( + f"user/key scope: Langfuse never received generation for key_alias={key_alias!r}" + ) + meta_blob = _json_blob(obs.metadata) + assert user_id in meta_blob or key_alias in (obs.name or ""), ( + f"user/key scope should attribute the user or key; metadata={meta_blob[:300]}" + ) + _assert_logs_spend( + client, + key=key, + outcome=outcome, + obs_cost=observation_spend(obs), + scope="user-key", + ) + + @pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["chat_completions"]) + def test_tool_calls_logged_with_cost( + self, + client: LoggingClient, + resources: ResourceManager, + langfuse_creds: LangfuseCreds, + ) -> None: + _, key, key_alias = self._user_key( + client, resources, langfuse_creds, models=[DRIVER_MODEL] + ) + prompt_marker = unique_marker() + outcome = client.chat_raw( + key, + DRIVER_MODEL, + f"Use get_weather for Tokyo. marker={prompt_marker}", + tools=[WEATHER_TOOL], + tool_choice="required", + max_tokens=128, + ) + require_successful_call(outcome) + + obs = client.poll_langfuse_observation( + langfuse_creds, + key_alias=key_alias, + prompt_marker=prompt_marker, + require_positive_cost=True, + ) + assert obs is not None + assert observation_mentions_tool(obs, "get_weather"), ( + f"user/key tool path: tool missing from Langfuse; output={str(obs.output)[:200]}" + ) + _assert_logs_spend( + client, + key=key, + outcome=outcome, + obs_cost=observation_spend(obs), + scope="user-key-tools", + ) + + +class TestLangfuseOrgScopedLogging: + """Org-scoped run: organization + team under it + team Langfuse callback. + + There is no /organization/.../callback today; logging attaches at the team + (or key) under the org. This class proves org-linked team keys still deliver + accurate Langfuse spend and team attribution (StandardLogging metadata + user_api_key_team_id / user_api_key_org_id). + """ + + def _org_team_key( + self, + client: LoggingClient, + resources: ResourceManager, + creds: LangfuseCreds, + *, + models: list[str], + ) -> tuple[str, str, str, str]: + marker = unique_marker() + key_alias = f"e2e-lf-org-key-{marker}" + org_id = client.create_org(f"e2e-lf-org-{marker}", models=models) + resources.defer(lambda: client.delete_org(org_id)) + team_id = client.create_team( + f"e2e-lf-org-team-{marker}", + models=models, + organization_id=org_id, + ) + resources.defer(lambda: client.delete_team(team_id)) + client.add_team_langfuse_callback(team_id, creds) + key = client.key_with_alias( + key_alias, + models=models, + team_id=team_id, + organization_id=org_id, + ) + resources.defer(lambda: client.delete_key(key)) + return org_id, team_id, key, key_alias + + @pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["chat_completions"]) + def test_success_logs_spend_with_team_attribution( + self, + client: LoggingClient, + resources: ResourceManager, + langfuse_creds: LangfuseCreds, + ) -> None: + org_id, team_id, key, key_alias = self._org_team_key( + client, resources, langfuse_creds, models=[DRIVER_MODEL] + ) + prompt_marker = unique_marker() + outcome = client.chat_raw( + key, DRIVER_MODEL, f"reply with one word only {prompt_marker}" + ) + require_successful_call(outcome) + + obs = client.poll_langfuse_observation( + langfuse_creds, + key_alias=key_alias, + prompt_marker=prompt_marker, + require_positive_cost=True, + ) + assert obs is not None, ( + f"org scope: Langfuse never received generation for key_alias={key_alias!r}" + ) + meta_blob = _json_blob(obs.metadata) + assert team_id in meta_blob, ( + f"org-scoped team key must stamp team_id on Langfuse metadata; " + f"team_id={team_id!r} metadata={meta_blob[:400]}" + ) + _ = org_id + _assert_logs_spend( + client, + key=key, + outcome=outcome, + obs_cost=observation_spend(obs), + scope="org-team", + ) diff --git a/tests/e2e/models.py b/tests/e2e/models.py index ab2835d87c4..e79d19215bf 100644 --- a/tests/e2e/models.py +++ b/tests/e2e/models.py @@ -23,6 +23,22 @@ class BudgetWindow(BaseModel): max_budget: float +class KeyLoggingCallbackVars(BaseModel): + langfuse_public_key: str | None = None + langfuse_secret_key: str | None = None + langfuse_host: str | None = None + + +class KeyLoggingCallback(BaseModel): + callback_name: str + callback_type: str = "success_and_failure" + callback_vars: KeyLoggingCallbackVars + + +class KeyMetadata(BaseModel): + logging: list[KeyLoggingCallback] | None = None + + class KeyGenerateBody(BaseModel): models: list[str] = [] duration: str | None = None @@ -31,6 +47,7 @@ class KeyGenerateBody(BaseModel): budget_duration: str | None = None user_id: str | None = None team_id: str | None = None + organization_id: str | None = None budget_id: str | None = None key_alias: str | None = None model_max_budget: dict[str, ModelBudgetEntry] | None = None @@ -39,6 +56,7 @@ class KeyGenerateBody(BaseModel): tpm_limit: int | None = None rpm_limit: int | None = None allowed_routes: list[str] | None = None + metadata: KeyMetadata | None = None class KeyGenerateResponse(BaseModel): @@ -105,6 +123,17 @@ class ThinkingParam(BaseModel): budget_tokens: int | None = None +class ChatToolFunction(BaseModel): + name: str + description: str | None = None + parameters: dict[str, object] | None = None + + +class ChatTool(BaseModel): + type: str = "function" + function: ChatToolFunction + + class ChatBody(BaseModel): model: str messages: list[ChatMessage] @@ -115,6 +144,9 @@ class ChatBody(BaseModel): reasoning_effort: str | None = None thinking: ThinkingParam | None = None service_tier: str | None = None + tools: list[ChatTool] | None = None + tool_choice: str | None = None + guardrails: list[str] | None = None class AnthropicMessagesBody(BaseModel): @@ -459,6 +491,7 @@ class TeamNewBody(BaseModel): team_alias: str models: list[str] = [] team_id: str | None = None + organization_id: str | None = None class TeamNewResponse(BaseModel):