test(e2e): cover Langfuse logging.yaml P0 logs_spend cells (#32857)

* test(e2e): cover Langfuse logging.yaml P0 logs_spend cells

Team, user/key, and org-scoped dynamic Langfuse callbacks drive real chat
traffic and assert calculatedTotalCost matches StandardLogging response_cost
and proxy spend. Also assert tool calls and applied guardrails land on the
trace. Missing env or proxy is a hard failure, never a skip

* test(e2e): use langfuse_otel callback for Langfuse spend coverage

Team and key dynamic logging attach callback_name=langfuse_otel (OTLP to
Langfuse) instead of the classic langfuse SDK. Match generations named
litellm_request by prompt marker and user_api_key_alias

* test(e2e): require Langfuse spend assert; drop AGENTS.md

Guardrail path no longer soft-gates logs_spend. Non-stream responses must
return positive x-litellm-response-cost; remove tests/e2e/AGENTS.md

* test(e2e): fail when Langfuse spend is missing on guardrail path

Always run logs_spend assertions for tool_permission; require positive
x-litellm-response-cost on non-stream and positive /spend/logs spend

* test(e2e): do not fall back to unmatched spend log rows

poll_proxy_spend_for_key returns None when response_id or positive-spend
filters match nothing, instead of silently using rows[0]
This commit is contained in:
mubashir1osmani 2026-07-11 14:11:00 -04:00 • committed by GitHub
parent 0bf81e2496
commit 1bf98c0687
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
6 changed files with 1164 additions and 53 deletions

View file

@ -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

View file

@ -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="<streamed>",
chunks=chunks,

View file

@ -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()

View file

@ -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 == "<streamed>":
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())

View file

@ -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",
)

View file

@ -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):