test: add logging e2e coverage (s3_v2, gcs_bucket, team langfuse callback, datadog failure)

Five new live e2e scenarios raising Logging & Guardrails registry coverage:
s3_v2 success and failure objects read back from the real S3 bucket,
gcs_bucket success record read back through the GCS JSON API (with
nextPageToken pagination and per-request bearer minting), team-scoped
Langfuse callback delivery with non-team isolation, and DataDog failure
event delivery queried by indexed model_group. datadog_reader gains
query-based variants of the marker search; the langfuse cell is a new
registry row. Bucket readers settle past a full flush interval so a
late duplicate cannot hide from the exactly-one assertions
This commit is contained in:
Yucheng Zhu 2026-08-27 10:33:06 -07:00
parent cd63c7e5a7
commit bcc2e71b5d
8 changed files with 835 additions and 33 deletions

View file

@ -24,3 +24,4 @@
- {id: logging.focus.success.writes_object, module: logging, tier: P1, event: success, assertions: [writes_object], exercised_on: [chat_completions, messages], source: "integrations/focus/focus_logger.py", rationale: "Cost mgmt multi-destination export"}
- {id: logging.niche_integrations.success.logs_spend, module: logging, tier: P2, event: success, assertions: [logs_spend], exercised_on: [chat_completions], source: grammar, rationale: "SMOKE cohort: athina/galileo/deepeval/langtrace/weave/lunary/humanloop/traceloop/helicone/argilla/newrelic/sqs/supabase/dynamodb/agentops/lago/etc"}
- {id: logging.niche_integrations.failure.logs_spend, module: logging, tier: P2, event: failure, assertions: [logs_spend], exercised_on: [chat_completions], source: grammar, rationale: "SMOKE niche failure path"}
- {id: logging.langfuse.success.logs_spend, module: logging, tier: P1, event: success, assertions: [logs_spend], exercised_on: [chat_completions], source: "integrations/langfuse/langfuse_otel.py", rationale: "Team-scoped Langfuse delivery via /team/callback; LangChain-ecosystem evals spend"}

View file

@ -97,15 +97,22 @@ class DdLogsReader:
indexed ``message`` empty, so a plain full-text query matches nothing;
``*:`` extends the scan to every attribute (the marker sits in the
prompt, e.g. ``messages.content``, wherever the route's payload puts
it). More than one hit for one call IS the duplicate-delivery bug, so
this never collapses to a single event. A 429 backs off and retries -
the search budget is org-wide, so another consumer can empty it under
us - while any other failure stays a hard fail."""
it)."""
return self.events_for_query(f"*:*{marker}*")
def events_for_query(self, query: str) -> list[DdLogEvent]:
"""Every ingested event the search query matches (failure payloads
carry no prompt to mark, so failure scenarios query indexed attributes
like ``@model_group:...`` instead of a body marker). More than one hit
for one call IS the duplicate-delivery bug, so this never collapses to
a single event. A 429 backs off and retries - the search budget is
org-wide, so another consumer can empty it under us - while any other
failure stays a hard fail."""
for _ in range(_RATE_LIMIT_RETRIES):
result = post(
URL(f"https://api.{self.site}/api/v2/logs/events/search"),
headers=_DdAuthHeaders(api_key=self.api_key, app_key=self.app_key),
json=_SearchRequest(filter=_SearchFilter(query=f"*:*{marker}*")),
json=_SearchRequest(filter=_SearchFilter(query=query)),
response_type=_SearchResponse,
timeout=30.0,
)
@ -123,6 +130,10 @@ class DdLogsReader:
)
def poll_events_for_marker(self, marker: str) -> list[DdLogEvent]:
"""``poll_events_for_query`` over the every-attribute marker scan."""
return self.poll_events_for_query(f"*:*{marker}*")
def poll_events_for_query(self, query: str) -> list[DdLogEvent]:
"""Poll until at least one matching event is searchable (the callback
flushes in periodic batches and DataDog ingestion adds seconds of lag),
then keep re-reading for DD_SETTLE_SECONDS so a late duplicate cannot
@ -132,15 +143,13 @@ class DdLogsReader:
request budget. At the deadline the last result is returned as-is."""
deadline = time.monotonic() + POLL_TIMEOUT
while time.monotonic() < deadline:
events = self.events_for_marker(marker)
events = self.events_for_query(query)
if events:
return self._settled_events_for_marker(marker, events)
return self._settled_events_for_query(query, events)
time.sleep(DD_SEARCH_INTERVAL)
return self.events_for_marker(marker)
return self.events_for_query(query)
def _settled_events_for_marker(
self, marker: str, events: list[DdLogEvent]
) -> list[DdLogEvent]:
def _settled_events_for_query(self, query: str, events: list[DdLogEvent]) -> list[DdLogEvent]:
"""Re-read at every search interval until the settle window closes; a
duplicate ends the watch early because more waiting cannot clear it.
@ -151,7 +160,7 @@ class DdLogsReader:
last_nonempty = events
while time.monotonic() < settle_deadline:
time.sleep(DD_SEARCH_INTERVAL)
latest = self.events_for_marker(marker)
latest = self.events_for_query(query)
if not latest:
continue
if len(latest) > 1:

View file

@ -0,0 +1,220 @@
"""Read-back for the gcs_bucket logging test against the real GCS bucket.
The proxy ships StandardLoggingPayload objects with its own service account
(litellm_settings.callbacks: ["gcs_bucket"] + GCS_BUCKET_NAME), and the test
reads them back through the GCS JSON API. Auth is a self-signed service-account
JWT (RS256 via PyJWT + cryptography, both litellm proxy dependencies the
runner installs) minted per request and sent directly as the Bearer token -
Google accepts that for storage.googleapis.com with no token exchange, which
keeps every HTTP read inside ``e2e_http``.
The default gcs_bucket mode batches payloads into ``{date}/batch-{id}.ndjson``
objects; unbatched mode writes ``{date}/{response_id}`` per call. The reader
handles both: it polls the day's listing, downloads the direct object when
present, and otherwise scans batch objects fresh enough to hold the call.
Missing configuration is a hard failure, never a skip.
"""
from __future__ import annotations
import os
import time
from dataclasses import dataclass
from datetime import datetime, timedelta, timezone
from pathlib import Path
from urllib.parse import quote
import jwt
import pytest
from pydantic import BaseModel, ConfigDict, Field
from e2e_config import POLL_INTERVAL, POLL_TIMEOUT
from e2e_http import URL, Headers, probe
_GCS_API = "https://storage.googleapis.com"
#: Tolerance for clock skew between this host and GCS object timestamps.
_SKEW = timedelta(seconds=120)
#: How long to keep re-reading after the first match before trusting the
#: exactly-one assertion: past one full gcs_bucket flush interval (~20s), so
#: a duplicate shipped by a later flush is seen, plus listing-latency margin.
GCS_SETTLE_SECONDS = 45.0
class _ServiceAccount(BaseModel):
model_config = ConfigDict(extra="ignore")
client_email: str
private_key: str
class _GcsAuthHeaders(Headers):
authorization: str = Field(serialization_alias="Authorization")
class _GcsObject(BaseModel):
model_config = ConfigDict(extra="ignore")
name: str
updated: datetime | None = None
class _GcsListResponse(BaseModel):
model_config = ConfigDict(extra="ignore")
items: list[_GcsObject] = []
next_page_token: str | None = Field(default=None, validation_alias="nextPageToken")
class _GcsListParams(BaseModel):
prefix: str
max_results: int = Field(default=1000, serialization_alias="maxResults")
page_token: str | None = Field(default=None, serialization_alias="pageToken")
class _GcsMediaParams(BaseModel):
alt: str = "media"
class GcsLogRecord(BaseModel):
"""The StandardLoggingPayload fields the gcs scenario pins."""
model_config = ConfigDict(extra="ignore")
id: str
status: str
model_group: str | None = None
response_cost: float | None = None
total_tokens: int | None = None
error_str: str | None = None
def _mint_bearer(account: _ServiceAccount) -> str:
"""Self-signed service-account JWT: for Google APIs a token whose ``aud``
is the service endpoint authorizes directly, no oauth2 token exchange.
Minted per request so a long session never outlives one token's expiry."""
now = int(time.time())
claims: dict[str, str | int] = {
"iss": account.client_email,
"sub": account.client_email,
"aud": f"{_GCS_API}/",
"iat": now,
"exp": now + 3600,
}
return jwt.encode(claims, account.private_key, algorithm="RS256")
@dataclass(frozen=True, slots=True)
class GcsLogReader:
bucket: str
account: _ServiceAccount
def _headers(self) -> _GcsAuthHeaders:
return _GcsAuthHeaders(authorization=f"Bearer {_mint_bearer(self.account)}")
def _list(self, prefix: str) -> list[_GcsObject]:
"""Every object under ``prefix``, following ``nextPageToken`` - the
shared day prefix accumulates all of the proxy's traffic, and a fresh
record past the 1000-object page cap must still be seen."""
items: list[_GcsObject] = []
page_token: str | None = None
while True:
result = probe(
URL(f"{_GCS_API}/storage/v1/b/{self.bucket}/o"),
headers=self._headers(),
params=_GcsListParams(prefix=prefix, page_token=page_token),
)
if result.status_code != 200:
pytest.fail(
f"GCS object listing for gs://{self.bucket}/{prefix} failed "
f"({result.status_code}): {result.body[:300]}"
)
page = _GcsListResponse.model_validate_json(result.body)
items.extend(page.items)
page_token = page.next_page_token
if not page_token:
return items
def _download(self, name: str) -> str:
result = probe(
URL(f"{_GCS_API}/storage/v1/b/{self.bucket}/o/{quote(name, safe='')}"),
headers=self._headers(),
params=_GcsMediaParams(),
)
if result.status_code != 200:
pytest.fail(
f"GCS object download gs://{self.bucket}/{name} failed ({result.status_code}): {result.body[:300]}"
)
return result.body
def records_for_response_id(self, response_id: str, *, since: datetime) -> list[GcsLogRecord]:
"""Every payload written for ``response_id``: the direct
``{date}/{response_id}`` object plus any hit inside batch NDJSON
objects updated after ``since``. More than one hit is the
duplicate-delivery bug, so this never collapses to a single record."""
records: list[GcsLogRecord] = []
window_start = since - _SKEW
for day_offset in (0, 1):
day = (since + timedelta(days=day_offset)).strftime("%Y-%m-%d")
for obj in self._list(f"{day}/"):
if obj.name == f"{day}/{response_id}":
records.append(GcsLogRecord.model_validate_json(self._download(obj.name)))
continue
is_fresh_batch = f"{day}/batch-" in obj.name and obj.updated is not None and obj.updated >= window_start
if is_fresh_batch:
records.extend(
GcsLogRecord.model_validate_json(line)
for line in self._download(obj.name).splitlines()
if response_id in line
)
return records
def poll_records_for_response_id(self, response_id: str, *, since: datetime) -> list[GcsLogRecord]:
"""Poll until the payload is readable (the gcs_bucket callback flushes
on a ~20s timer), then keep re-reading for GCS_SETTLE_SECONDS - past a
full flush interval - so a duplicate shipped by a later flush cannot
hide from the exactly-one assertion. A duplicate ends the settle early
because more waiting cannot clear it."""
deadline = time.monotonic() + POLL_TIMEOUT
while time.monotonic() < deadline:
records = self.records_for_response_id(response_id, since=since)
if records:
return self._settled_records(response_id, since=since, first=records)
time.sleep(POLL_INTERVAL)
return []
def _settled_records(self, response_id: str, *, since: datetime, first: list[GcsLogRecord]) -> list[GcsLogRecord]:
"""Re-read at every poll interval until the settle window closes; a
transiently empty re-read never downgrades what was already seen."""
settle_deadline = time.monotonic() + GCS_SETTLE_SECONDS
latest = first
while time.monotonic() < settle_deadline and len(latest) <= 1:
time.sleep(POLL_INTERVAL)
latest = self.records_for_response_id(response_id, since=since) or latest
return latest
def utc_now() -> datetime:
return datetime.now(timezone.utc)
def build_gcs_reader() -> GcsLogReader:
bucket = os.environ.get("GCS_BUCKET_NAME", "")
if not bucket:
pytest.fail(
"GCS_BUCKET_NAME must be set: the gcs test reads the proxy's gcs_bucket "
"delivery back from the real bucket (the cluster secret manager injects "
"it; locally set it in tests/e2e/.env)"
)
raw = ""
credentials_path = os.environ.get("GOOGLE_APPLICATION_CREDENTIALS", "")
if credentials_path and Path(credentials_path).is_file():
raw = Path(credentials_path).read_text()
else:
raw = os.environ.get("VERTEXAI_CREDENTIALS", "")
if not raw:
pytest.fail(
"GCS read-back needs a service-account key: set "
"GOOGLE_APPLICATION_CREDENTIALS (path) or VERTEXAI_CREDENTIALS (JSON), "
"as the cluster secret manager does"
)
return GcsLogReader(bucket=bucket, account=_ServiceAccount.model_validate_json(raw))

View file

@ -0,0 +1,115 @@
"""Read-back for the s3 logging tests against the real S3 bucket the proxy
ships StandardLoggingPayload objects to (litellm_settings.callbacks: ["s3_v2"]).
Delivery is judged on what actually landed in the bucket: the proxy writes
with its own credentials exactly as in production, and the tests list and
download the objects back with boto3 (already a litellm proxy dependency, so
the e2e runner image carries it; it is an AWS SDK, not a raw HTTP client, so
the e2e_http-only transport rule is untouched). The bucket comes from
S3_LOGS_BUCKET_NAME - on the cluster the secret manager injects it, locally
tests/e2e/.env provides it. Missing configuration is a hard failure, never a
skip.
"""
from __future__ import annotations
import os
import time
from collections.abc import Callable
from dataclasses import dataclass
from typing import TYPE_CHECKING
import boto3
import pytest
from pydantic import BaseModel, ConfigDict
from e2e_config import POLL_INTERVAL, POLL_TIMEOUT
if TYPE_CHECKING:
from types_boto3_s3.client import S3Client
#: How long to keep re-reading after the first match before trusting the
#: exactly-one assertion: past one full s3_v2 flush interval (~10s), so a
#: duplicate shipped by a LATER flush is seen, plus listing-latency margin.
#: The DataDog reader settles the same way (DD_SETTLE_SECONDS).
S3_SETTLE_SECONDS = 25.0
class S3LogRecord(BaseModel):
"""The StandardLoggingPayload fields the s3 scenarios pin."""
model_config = ConfigDict(extra="ignore")
id: str
status: str
model_group: str | None = None
response_cost: float | None = None
total_tokens: int | None = None
error_str: str | None = None
@dataclass(frozen=True, slots=True)
class S3LogReader:
bucket: str
client: S3Client
def list_keys(self, prefix: str) -> list[str]:
response = self.client.list_objects_v2(Bucket=self.bucket, Prefix=prefix)
return [obj["Key"] for obj in response.get("Contents", []) if "Key" in obj]
def read_record(self, key: str) -> S3LogRecord:
body = self.client.get_object(Bucket=self.bucket, Key=key)["Body"].read()
return S3LogRecord.model_validate_json(body)
def records_matching(self, *, prefix: str, predicate: Callable[[S3LogRecord], bool]) -> list[S3LogRecord]:
return [record for record in map(self.read_record, self.list_keys(prefix)) if predicate(record)]
def poll_records(self, *, prefix: str, predicate: Callable[[S3LogRecord], bool]) -> list[S3LogRecord]:
"""Poll until at least one matching object is listed (the s3_v2
callback flushes on a ~10s timer), then keep re-reading for
S3_SETTLE_SECONDS - past a full flush interval - so a duplicate
shipped by a later flush cannot hide from the exactly-one assertion.
One blind spot is inherent: a duplicate write that reuses the exact
same object key overwrites the first object and no listing can see
it; distinct-key duplicates are what this catches. At the deadline an
empty list is returned and the caller's assertion carries the failure
message."""
deadline = time.monotonic() + POLL_TIMEOUT
while time.monotonic() < deadline:
records = self.records_matching(prefix=prefix, predicate=predicate)
if records:
return self._settled_records(prefix=prefix, predicate=predicate, first=records)
time.sleep(POLL_INTERVAL)
return []
def _settled_records(
self, *, prefix: str, predicate: Callable[[S3LogRecord], bool], first: list[S3LogRecord]
) -> list[S3LogRecord]:
"""Re-read at every poll interval until the settle window closes; a
duplicate ends the watch early because more waiting cannot clear it.
A transiently empty re-read never downgrades what was already seen."""
settle_deadline = time.monotonic() + S3_SETTLE_SECONDS
latest = first
while time.monotonic() < settle_deadline and len(latest) <= 1:
time.sleep(POLL_INTERVAL)
latest = self.records_matching(prefix=prefix, predicate=predicate) or latest
return latest
def build_s3_reader() -> S3LogReader:
bucket = os.environ.get("S3_LOGS_BUCKET_NAME", "")
if not bucket:
pytest.fail(
"S3_LOGS_BUCKET_NAME must be set: the s3 tests read the proxy's s3_v2 "
"delivery back from the real bucket (the cluster secret manager injects "
"it; locally set it in tests/e2e/.env to the same bucket "
"s3_callback_params.s3_bucket_name names)"
)
region = os.environ.get("AWS_REGION_NAME") or os.environ.get("AWS_REGION") or "us-east-1"
return S3LogReader(
bucket=bucket,
# boto3.client's overload set covers every AWS service; the ones without
# installed stubs type as Unknown, so the member is "partially unknown"
# even though the s3 overload itself resolves to S3Client.
client=boto3.client("s3", region_name=region), # pyright: ignore[reportUnknownMemberType]
)

View file

@ -19,6 +19,7 @@ received).
from __future__ import annotations
import math
import time
import pytest
from pydantic import BaseModel, ConfigDict
@ -27,7 +28,8 @@ from datadog_reader import DdLogEvent, DdLogsReader
from e2e_config import CHEAP_ANTHROPIC_MODEL, CHEAP_OPENAI_MODEL, unique_marker
from e2e_http import NoBody
from lifecycle import ResourceManager
from logging_client import LoggingClient, first_ok
from logging_client import INVALID_UPSTREAM_API_KEY, LoggingClient, first_ok
from models import LiteLLMParamsBody
pytestmark = pytest.mark.e2e
@ -46,6 +48,7 @@ class _DdMessagePayload(BaseModel):
status: str
call_type: str
stream: bool | None = None
error_str: str | None = None
def _assert_datadog_configured(client: LoggingClient) -> None:
@ -89,18 +92,14 @@ def _assert_exactly_one_event(
# indexed event status from the parsed payload's status attribute
# ("success") and normalizes it to its OK severity - so "ok" is what a
# successfully ingested success event looks like on the search API.
assert event.status == "ok", (
f"success events must index at DataDog's ok severity, got {event.status!r}"
)
assert event.status == "ok", f"success events must index at DataDog's ok severity, got {event.status!r}"
payload = _DdMessagePayload.model_validate(event.attributes)
assert payload.status == "success", f"payload status must be success, got {payload.status!r}"
assert payload.model_group == model_group, (
f"payload model_group must be {model_group!r}, got {payload.model_group!r}"
)
assert payload.call_type == call_type, (
f"payload call_type must be {call_type!r}, got {payload.call_type!r}"
)
assert payload.call_type == call_type, f"payload call_type must be {call_type!r}, got {payload.call_type!r}"
assert payload.total_tokens > 0, f"payload must count real tokens, got {payload.total_tokens}"
# Relative tolerance, not bit-equality: the cost round-trips through
# DataDog's attribute indexing, whose float serialization may drift in the
@ -109,9 +108,7 @@ def _assert_exactly_one_event(
f"payload response_cost {payload.response_cost} must equal the anchor cost {cost_anchor}"
)
if expect_stream:
assert payload.stream is True, (
f"a streamed call's payload must record stream=true, got {payload.stream!r}"
)
assert payload.stream is True, f"a streamed call's payload must record stream=true, got {payload.stream!r}"
return payload
@ -211,7 +208,9 @@ class TestDataDogLogDelivery:
marker = unique_marker()
outcome = first_ok(
client,
lambda: client.chat_raw(key, CHEAP_ANTHROPIC_MODEL, f"reply with one word {marker}", stream=True, max_tokens=16),
lambda: client.chat_raw(
key, CHEAP_ANTHROPIC_MODEL, f"reply with one word {marker}", stream=True, max_tokens=16
),
)
assert outcome.is_streaming, f"response must be an event stream, got content-type {outcome.content_type!r}"
assert outcome.chunks > 0, "the stream must deliver at least one event"
@ -231,9 +230,7 @@ class TestDataDogLogDelivery:
cost_anchor=spend_row.spend,
expect_stream=True,
)
assert spend_row.total_tokens is not None, (
"the spend row must record total_tokens for the token cross-check"
)
assert spend_row.total_tokens is not None, "the spend row must record total_tokens for the token cross-check"
assert spend_row.total_tokens == payload.total_tokens, (
f"the spend row and the DataDog event must agree on tokens: "
f"{spend_row.total_tokens} vs {payload.total_tokens}"
@ -255,7 +252,9 @@ class TestDataDogLogDelivery:
marker = unique_marker()
outcome = first_ok(
client,
lambda: client.messages_raw(key, CHEAP_ANTHROPIC_MODEL, f"reply with one word {marker}", max_tokens=16, stream=True),
lambda: client.messages_raw(
key, CHEAP_ANTHROPIC_MODEL, f"reply with one word {marker}", max_tokens=16, stream=True
),
)
assert outcome.is_streaming, f"response must be an event stream, got content-type {outcome.content_type!r}"
assert outcome.chunks > 0, "the stream must deliver at least one event"
@ -275,9 +274,7 @@ class TestDataDogLogDelivery:
cost_anchor=spend_row.spend,
expect_stream=True,
)
assert spend_row.total_tokens is not None, (
"the spend row must record total_tokens for the token cross-check"
)
assert spend_row.total_tokens is not None, "the spend row must record total_tokens for the token cross-check"
assert spend_row.total_tokens == payload.total_tokens, (
f"the spend row and the DataDog event must agree on tokens: "
f"{spend_row.total_tokens} vs {payload.total_tokens}"
@ -319,10 +316,89 @@ class TestDataDogLogDelivery:
cost_anchor=spend_row.spend,
expect_stream=True,
)
assert spend_row.total_tokens is not None, (
"the spend row must record total_tokens for the token cross-check"
)
assert spend_row.total_tokens is not None, "the spend row must record total_tokens for the token cross-check"
assert spend_row.total_tokens == payload.total_tokens, (
f"the spend row and the DataDog event must agree on tokens: "
f"{spend_row.total_tokens} vs {payload.total_tokens}"
)
def _assert_exactly_one_failure_event(events: list[DdLogEvent], *, model_group: str) -> _DdMessagePayload:
"""The enforced behavior for a failed call: the intake holds exactly one
event for the deployment, sourced from litellm, indexed at an error-grade
severity (DataDog derives it from the payload's status="failure"; observed
as its "emergency" bucket), whose payload carries the provider error and
no cost."""
assert events, "no DataDog log event for the failed call reached the intake within the deadline"
assert len(events) == 1, (
f"expected exactly ONE DataDog log event for the failed call, got {len(events)} - "
"more than one event for one call is the duplicate-delivery bug"
)
event = events[0]
assert "source:litellm" in event.tags, (
f"the ingested event must carry the litellm source (shipped as ddsource), got tags {event.tags!r}"
)
assert event.status in ("error", "emergency"), (
f"failure events must index at an error-grade severity, got {event.status!r}"
)
payload = _DdMessagePayload.model_validate(event.attributes)
assert payload.status == "failure", f"payload status must be failure, got {payload.status!r}"
assert payload.model_group == model_group, (
f"payload model_group must be {model_group!r}, got {payload.model_group!r}"
)
assert not payload.response_cost, f"a failed call must not be billed, got response_cost={payload.response_cost!r}"
return payload
class TestDataDogFailureDelivery:
@pytest.mark.covers("logging.datadog.failure.exports_metric", exercised_on=["chat_completions"])
def test_failed_chat_completions_emits_one_error_event(
self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager
) -> None:
"""A /chat/completions call that fails at the provider must reach the
DataDog logs intake as exactly one error-grade event carrying the
provider error - failure metrics drive alerting and SLOs, so a dropped
failure event is an invisible outage.
A deployment with an invalid upstream key lets the request pass proxy
auth and fail at the provider (the same lever as the OTEL error test).
Failure payloads carry no prompt to mark, so the read-back queries the
indexed @model_group attribute of the per-run unique deployment name;
proxy-side 401s during key propagation never reach the provider and
ship no payload, so exactly one provider failure exists for it."""
_assert_datadog_configured(client)
model_name = f"dd-err-{unique_marker()}"
model_id = client.create_model(
model_name,
LiteLLMParamsBody(model="anthropic/claude-haiku-4-5", api_key=INVALID_UPSTREAM_API_KEY),
)
resources.defer(lambda: client.delete_model(model_id))
key = client.key_with_alias(f"dd-err-key-{unique_marker()}", models=[model_name])
resources.defer(lambda: client.delete_key(key))
deadline = time.monotonic() + client.proxy.poll_timeout
while True:
outcome = client.chat_raw(key, model_name, "trigger an upstream auth failure", max_tokens=16)
assert not outcome.ok, "the call must fail; the deployment's upstream key is invalid"
assert outcome.status_code != -1, (
"network failure between the test and the proxy while provoking the provider "
"failure; retrying now could double-log the failure payload and falsely trip "
f"the exactly-one assertion - fix the rig connectivity first: {outcome.body[:200]}"
)
if "AnthropicException" in outcome.body or time.monotonic() >= deadline:
break
time.sleep(client.proxy.poll_interval)
assert "AnthropicException" in outcome.body, (
"never saw the upstream provider failure before the deadline; the key may still be "
f"propagating - last outcome {outcome.status_code}: {outcome.body[:200]}"
)
assert outcome.status_code == 401, (
f"an upstream auth failure must map to 401, got {outcome.status_code}: {outcome.body[:200]}"
)
events = dd_logs.poll_events_for_query(f"@model_group:{model_name}")
payload = _assert_exactly_one_failure_event(events, model_group=model_name)
assert payload.error_str is not None and "AnthropicException" in payload.error_str, (
f"the event must carry the provider error, got error_str={payload.error_str!r}"
)

View file

@ -0,0 +1,101 @@
"""Live e2e: gcs_bucket log delivery for successful calls.
Covers logging.gcs_bucket.success.writes_object: one successful
/chat/completions call must land in the real GCS bucket as exactly one
StandardLoggingPayload record (GCS is the audit-trail parallel to S3 for GCP
deployments). Delivery is judged on what is actually readable in the bucket:
the proxy writes with its production service account, and the test reads the
record back through the GCS JSON API - covering both the batched NDJSON layout
(the default) and the per-request object layout.
Both halves of the contract are asserted: the recorded state (the proxy
reports the GCSBucketLogger callback active via /health/readiness/details -
note gcs_bucket is enterprise-gated, so this also requires a license) and the
enforced behavior (the record in the bucket, cost cross-checked against the
x-litellm-response-cost header of the very response the caller received).
"""
from __future__ import annotations
import math
import pytest
from e2e_config import CHEAP_ANTHROPIC_MODEL, unique_marker
from e2e_http import NoBody
from gcs_reader import GcsLogReader, build_gcs_reader, utc_now
from lifecycle import ResourceManager
from logging_client import LoggingClient, completion_response_id, first_ok
pytestmark = pytest.mark.e2e
#: The active gcs_bucket callback's name in /health/readiness/details success_callbacks.
GCS_LOGGER_NAME = "GCSBucketLogger"
@pytest.fixture(scope="session")
def gcs_logs() -> GcsLogReader:
return build_gcs_reader()
def _assert_gcs_configured(client: LoggingClient) -> None:
"""Recorded state: the proxy reports the gcs_bucket callback among its
active callbacks, so a missing destination config (or a missing enterprise
license - gcs_bucket refuses to initialize without one) fails here, before
any delivery-based assertion can time out confusingly."""
result = client.proxy.probe("/health/readiness/details", params=NoBody())
assert result.status_code == 200, (
f"/health/readiness/details must answer 200, got {result.status_code}: {result.body[:300]}"
)
assert GCS_LOGGER_NAME in result.body, (
f"the proxy must report the {GCS_LOGGER_NAME} callback active "
f"(litellm_settings.callbacks: ['gcs_bucket'] + GCS_BUCKET_NAME env + enterprise license); "
f"got: {result.body[:400]}"
)
class TestGcsLogDelivery:
@pytest.mark.covers("logging.gcs_bucket.success.writes_object", exercised_on=["chat_completions"])
def test_chat_completions_writes_one_success_record(
self, client: LoggingClient, gcs_logs: GcsLogReader, resources: ResourceManager
) -> None:
"""One successful non-streaming /chat/completions call must be
readable back from the bucket as exactly one payload record carrying
the model group, the token counts, and the same cost the caller's
response header reported."""
_assert_gcs_configured(client)
alias = f"gcs-chat-{unique_marker()}"
key = client.key_with_alias(alias, models=[CHEAP_ANTHROPIC_MODEL])
resources.defer(lambda: client.delete_key(key))
since = utc_now()
marker = unique_marker()
outcome = first_ok(
client,
lambda: client.chat_raw(key, CHEAP_ANTHROPIC_MODEL, f"reply with one word {marker}", max_tokens=16),
)
assert outcome.response_cost is not None and outcome.response_cost > 0, (
f"the response must report x-litellm-response-cost, got {outcome.response_cost!r}"
)
body_id = completion_response_id(outcome.body)
assert body_id is not None, "the completion body must carry an id (it names the gcs record)"
records = gcs_logs.poll_records_for_response_id(body_id, since=since)
assert records, f"no gcs record for response {body_id} was readable from the bucket within the deadline"
assert len(records) == 1, (
f"expected exactly ONE gcs record for the call, got {len(records)} - "
"more than one record for one call is the duplicate-delivery bug"
)
record = records[0]
assert record.id == body_id, f"record id must be the response id, got {record.id!r}"
assert record.status == "success", f"payload status must be success, got {record.status!r}"
assert record.model_group == CHEAP_ANTHROPIC_MODEL, (
f"payload model_group must be {CHEAP_ANTHROPIC_MODEL!r}, got {record.model_group!r}"
)
assert record.total_tokens is not None and record.total_tokens > 0, (
f"payload must count real tokens, got {record.total_tokens!r}"
)
assert record.response_cost is not None and math.isclose(
record.response_cost, outcome.response_cost, rel_tol=1e-9
), f"payload response_cost {record.response_cost!r} must equal the header cost {outcome.response_cost}"

View file

@ -0,0 +1,168 @@
"""Live e2e: s3_v2 log delivery for successful and failed calls.
Covers logging.s3.success.writes_object and logging.s3.failure.writes_object:
one /chat/completions call must land in the real S3 bucket as exactly one
StandardLoggingPayload object (the primary audit trail; the batch flush must
neither drop nor duplicate it), and a failed call must be persisted the same
way for compliance. Delivery is judged on what is actually in the bucket: the
proxy writes with its production credentials and the test lists and reads the
objects back.
Both halves of the contract are asserted: the recorded state (the proxy
reports the S3Logger callback active via /health/readiness/details) and the
enforced behavior (the object in the bucket, with the cost cross-checked
against the x-litellm-response-cost header of the very response the caller
received).
The suite requires ``s3_callback_params.s3_use_key_prefix: true`` on the proxy,
which keys objects as ``{key_alias}/{date}/time-..._{id}.json`` - a unique key
alias per test turns the poll into a cheap prefix listing.
"""
from __future__ import annotations
import math
import time
import pytest
from e2e_config import CHEAP_ANTHROPIC_MODEL, unique_marker
from e2e_http import NoBody
from lifecycle import ResourceManager
from logging_client import (
INVALID_UPSTREAM_API_KEY,
LoggingClient,
completion_response_id,
first_ok,
)
from models import LiteLLMParamsBody
from s3_reader import S3LogReader, build_s3_reader
pytestmark = pytest.mark.e2e
#: The active s3_v2 callback's name in /health/readiness/details success_callbacks.
S3_LOGGER_NAME = "S3Logger"
@pytest.fixture(scope="session")
def s3_logs() -> S3LogReader:
return build_s3_reader()
def _assert_s3_configured(client: LoggingClient) -> None:
"""Recorded state: the proxy reports the s3_v2 callback among its active
callbacks, so a missing destination config fails here, before any
delivery-based assertion can time out confusingly."""
result = client.proxy.probe("/health/readiness/details", params=NoBody())
assert result.status_code == 200, (
f"/health/readiness/details must answer 200, got {result.status_code}: {result.body[:300]}"
)
assert S3_LOGGER_NAME in result.body, (
f"the proxy must report the {S3_LOGGER_NAME} callback active "
f"(litellm_settings.callbacks: ['s3_v2'] + s3_callback_params in the proxy config); "
f"got: {result.body[:400]}"
)
class TestS3LogDelivery:
@pytest.mark.covers("logging.s3.success.writes_object", exercised_on=["chat_completions"])
def test_chat_completions_writes_one_success_object(
self, client: LoggingClient, s3_logs: S3LogReader, resources: ResourceManager
) -> None:
"""One successful non-streaming /chat/completions call must land in
the bucket as exactly one payload object carrying the model group, the
token counts, and the same cost the caller's response header reported."""
_assert_s3_configured(client)
alias = f"s3-chat-{unique_marker()}"
key = client.key_with_alias(alias, models=[CHEAP_ANTHROPIC_MODEL])
resources.defer(lambda: client.delete_key(key))
marker = unique_marker()
outcome = first_ok(
client,
lambda: client.chat_raw(key, CHEAP_ANTHROPIC_MODEL, f"reply with one word {marker}", max_tokens=16),
)
assert outcome.response_cost is not None and outcome.response_cost > 0, (
f"the response must report x-litellm-response-cost, got {outcome.response_cost!r}"
)
body_id = completion_response_id(outcome.body)
assert body_id is not None, "the completion body must carry an id (it names the s3 object)"
records = s3_logs.poll_records(prefix=f"{alias}/", predicate=lambda r: r.id == body_id)
assert records, (
f"no s3 object for response {body_id} under prefix {alias}/ reached the bucket within the deadline"
)
assert len(records) == 1, (
f"expected exactly ONE s3 object for the call, got {len(records)} - "
"more than one object for one call is the duplicate-delivery bug"
)
record = records[0]
assert record.status == "success", f"payload status must be success, got {record.status!r}"
assert record.model_group == CHEAP_ANTHROPIC_MODEL, (
f"payload model_group must be {CHEAP_ANTHROPIC_MODEL!r}, got {record.model_group!r}"
)
assert record.total_tokens is not None and record.total_tokens > 0, (
f"payload must count real tokens, got {record.total_tokens!r}"
)
assert record.response_cost is not None and math.isclose(
record.response_cost, outcome.response_cost, rel_tol=1e-9
), f"payload response_cost {record.response_cost!r} must equal the header cost {outcome.response_cost}"
@pytest.mark.covers("logging.s3.failure.writes_object", exercised_on=["chat_completions"])
def test_chat_completions_failure_writes_one_object(
self, client: LoggingClient, s3_logs: S3LogReader, resources: ResourceManager
) -> None:
"""A call that fails at the provider must be persisted to the bucket as
exactly one failure payload carrying the provider error - failed calls
are part of the audit trail, not an exemption from it.
A deployment with an invalid upstream key lets the request pass proxy
auth and fail at the provider (the same lever as the OTEL error test);
proxy-side 401s during key propagation never reach the provider and
ship no payload, so exactly one provider failure exists for this alias."""
_assert_s3_configured(client)
model_name = f"s3-err-{unique_marker()}"
model_id = client.create_model(
model_name,
LiteLLMParamsBody(model="anthropic/claude-haiku-4-5", api_key=INVALID_UPSTREAM_API_KEY),
)
resources.defer(lambda: client.delete_model(model_id))
alias = f"s3-err-key-{unique_marker()}"
key = client.key_with_alias(alias, models=[model_name])
resources.defer(lambda: client.delete_key(key))
deadline = time.monotonic() + client.proxy.poll_timeout
while True:
outcome = client.chat_raw(key, model_name, "trigger an upstream auth failure", max_tokens=16)
assert not outcome.ok, "the call must fail; the deployment's upstream key is invalid"
assert outcome.status_code != -1, (
"network failure between the test and the proxy while provoking the provider "
"failure; retrying now could double-log the failure payload and falsely trip "
f"the exactly-one assertion - fix the rig connectivity first: {outcome.body[:200]}"
)
if "AnthropicException" in outcome.body or time.monotonic() >= deadline:
break
time.sleep(client.proxy.poll_interval)
assert "AnthropicException" in outcome.body, (
"never saw the upstream provider failure before the deadline; the key may still be "
f"propagating - last outcome {outcome.status_code}: {outcome.body[:200]}"
)
assert outcome.status_code == 401, (
f"an upstream auth failure must map to 401, got {outcome.status_code}: {outcome.body[:200]}"
)
records = s3_logs.poll_records(
prefix=f"{alias}/",
predicate=lambda r: r.status == "failure" and r.model_group == model_name,
)
assert records, (
f"no failure object for {model_name} under prefix {alias}/ reached the bucket within the deadline"
)
assert len(records) == 1, f"expected exactly ONE failure object for the call, got {len(records)}"
record = records[0]
assert record.error_str is not None and "AnthropicException" in record.error_str, (
f"the persisted failure must carry the provider error, got error_str={record.error_str!r}"
)
assert not record.response_cost, f"a failed call must not be billed, got response_cost={record.response_cost!r}"

View file

@ -0,0 +1,112 @@
"""Live e2e: team-scoped Langfuse callback delivery and isolation.
Covers logging.langfuse.success.logs_spend: a team configured with a Langfuse
callback via POST /team/{id}/callback must deliver its members' calls to the
real Langfuse project (generation readable back through Langfuse's own API,
with the cost agreeing with the x-litellm-response-cost header), while traffic
from keys outside the team must NOT reach that project - the isolation is the
point of team-scoped callbacks.
Both halves of the contract are asserted: the recorded state (the /team/callback
registration itself answers success) and the enforced behavior (the generation
at the destination for the team key, and its absence for the non-team key).
"""
from __future__ import annotations
import time
import pytest
from e2e_config import CHEAP_ANTHROPIC_MODEL, unique_marker
from lifecycle import ResourceManager
from logging_client import (
LangfuseCreds,
LoggingClient,
costs_agree,
first_ok,
load_langfuse_creds,
observation_spend,
)
pytestmark = pytest.mark.e2e
#: How long to keep re-checking that the non-team call never surfaces in
#: Langfuse after the team call's generation has already been ingested; the
#: positive observation bounds the pipeline's latency, so a wrong delivery
#: would be visible within the same order of magnitude.
ISOLATION_SETTLE_SECONDS = 30.0
ISOLATION_CHECK_INTERVAL_SECONDS = 5.0
@pytest.fixture(scope="session")
def langfuse_creds() -> LangfuseCreds:
return load_langfuse_creds()
class TestTeamLangfuseCallback:
@pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["chat_completions"])
def test_team_callback_delivers_and_isolates(
self, client: LoggingClient, langfuse_creds: LangfuseCreds, resources: ResourceManager
) -> None:
team_id = client.create_team(f"lf-team-{unique_marker()}", models=[CHEAP_ANTHROPIC_MODEL])
resources.defer(lambda: client.delete_team(team_id))
# Recorded state: the registration endpoint itself must answer success
# (add_team_langfuse_callback asserts it).
client.add_team_langfuse_callback(team_id, langfuse_creds)
team_alias = f"lf-team-key-{unique_marker()}"
team_key = client.key_with_alias(team_alias, models=[CHEAP_ANTHROPIC_MODEL], team_id=team_id)
resources.defer(lambda: client.delete_key(team_key))
solo_alias = f"lf-solo-key-{unique_marker()}"
solo_key = client.key_with_alias(solo_alias, models=[CHEAP_ANTHROPIC_MODEL])
resources.defer(lambda: client.delete_key(solo_key))
team_marker = unique_marker()
team_outcome = first_ok(
client,
lambda: client.chat_raw(
team_key, CHEAP_ANTHROPIC_MODEL, f"reply with one word {team_marker}", max_tokens=16
),
)
assert team_outcome.response_cost is not None and team_outcome.response_cost > 0, (
f"the response must report x-litellm-response-cost, got {team_outcome.response_cost!r}"
)
solo_marker = unique_marker()
_ = first_ok(
client,
lambda: client.chat_raw(
solo_key, CHEAP_ANTHROPIC_MODEL, f"reply with one word {solo_marker}", max_tokens=16
),
)
# Enforced behavior, positive half: the team member's call is readable
# back from the real Langfuse project with an agreeing cost.
observation = client.poll_langfuse_observation(
langfuse_creds,
key_alias=team_alias,
prompt_marker=team_marker,
require_positive_cost=True,
)
assert observation is not None, (
f"the team key's call (marker {team_marker}) never reached Langfuse within the deadline"
)
cost = observation_spend(observation)
assert cost is not None and costs_agree(team_outcome.response_cost, cost), (
f"Langfuse calculatedTotalCost {cost!r} must agree with the header cost {team_outcome.response_cost}"
)
# Enforced behavior, negative half: the non-team call must never show
# up in this project. The positive generation above has already been
# ingested, which bounds the pipeline latency, so keep re-checking for
# a settle window rather than trusting a single instant.
settle_deadline = time.monotonic() + ISOLATION_SETTLE_SECONDS
while True:
leaked = client.find_langfuse_observation(langfuse_creds, key_alias=solo_alias, prompt_marker=solo_marker)
assert leaked is None, (
f"a non-team key's call (marker {solo_marker}) reached the team's Langfuse "
f"project: {leaked.id} - team callbacks must not apply outside the team"
)
if time.monotonic() >= settle_deadline:
break
time.sleep(ISOLATION_CHECK_INTERVAL_SECONDS)