merge(e2e): PR #34650 into combined e2e run branch

This commit is contained in:
mubashir1osmani 2026-08-10 23:00:57 -07:00
commit 05e4559352
7 changed files with 253 additions and 1 deletions

View file

@ -34,6 +34,8 @@
- {id: llm.rerank.cohere.basic.nonstream.works, module: llm, tier: P1, subject_endpoint: rerank, route: cohere, capability: basic, streaming: nonstream, assertions: [works], source: "test_rerank_e2e.py:29", rationale: "Cohere rerank, top_n + relevance_score"}
- {id: llm.files.openai.content.nonstream.works, module: llm, tier: P0, subject_endpoint: files, route: openai, capability: basic, streaming: nonstream, assertions: [works], source: "test_batches_e2e.py", rationale: "GET /v1/files/{id}/content returns uploaded batch JSONL bytes"}
- {id: llm.realtime.bedrock_converse.basic.stream.works, module: llm, tier: P0, subject_endpoint: realtime, route: bedrock_converse, capability: basic, streaming: stream, assertions: [works], source: "test_realtime_bedrock_e2e.py", rationale: "Nova Sonic realtime session emits response.done (LIT-2239)"}
- {id: llm.google_native.gemini.basic.nonstream.cost_logged, module: llm, tier: P0, subject_endpoint: google_native, route: gemini, capability: basic, streaming: nonstream, assertions: [cost_logged], source: "LIT-4076 / proxy/google_endpoints/endpoints.py", fail_before_fix: proven, rationale: "google-native generateContent must stamp x-litellm-response-cost so SDK traffic reconciles against spend"}
- {id: llm.google_native.gemini.basic.stream.works, module: llm, tier: P0, subject_endpoint: google_native, route: gemini, capability: basic, streaming: stream, assertions: [works], source: "PR #28213 / proxy/proxy_server.py async_data_generator", fail_before_fix: proven, rationale: "streamGenerateContent must relay single-prefixed SSE frames with no [DONE] sentinel; doubled data: prefixes and the OpenAI terminator both break the Vertex Java SDK"}
- {id: llm.rerank.bedrock.basic.nonstream.works, module: llm, tier: P1, subject_endpoint: rerank, route: bedrock_converse, capability: basic, streaming: nonstream, assertions: [works], source: "llms/bedrock/rerank/handler.py", rationale: "Bedrock rerank"}
- {id: llm.rerank.together_ai.basic.nonstream.works, module: llm, tier: P1, subject_endpoint: rerank, route: together_ai, capability: basic, streaming: nonstream, assertions: [works], source: "llms/together_ai/rerank/handler.py", rationale: "Together rerank"}
- {id: llm.images_generations.openai.basic.nonstream.works, module: llm, tier: P1, subject_endpoint: images_generations, route: openai, capability: basic, streaming: nonstream, assertions: [works], source: "test_image_generation_e2e.py:22", rationale: "OpenAI image gen, b64/url"}

View file

@ -6,6 +6,7 @@
- {id: logging.datadog.stream.exports_metric, module: logging, tier: P0, event: stream, assertions: [exports_metric], exercised_on: [chat_completions, messages, responses], source: "integrations/datadog/datadog.py", rationale: "Streaming aggregates usage after the last chunk; delivery and cost must survive that path"}
- {id: logging.datadog.failure.exports_metric, module: logging, tier: P0, event: failure, assertions: [exports_metric], exercised_on: [chat_completions], source: "integrations/datadog/datadog.py", rationale: "Failure metrics for alerting/SLO"}
- {id: logging.prometheus.success.exports_metric, module: logging, tier: P0, event: success, assertions: [exports_metric], exercised_on: [chat_completions, messages, embeddings], source: "integrations/prometheus.py", rationale: "Standard OSS metrics; per-key cardinality (existing e2e)"}
- {id: logging.prometheus.success.records_queue_time, module: logging, tier: P1, event: success, assertions: [records_queue_time], exercised_on: [chat_completions], source: "integrations/prometheus.py / LIT-2034", fail_before_fix: proven, rationale: "Queue time feeds saturation alerting; the family stayed registered while no observation was ever recorded, so presence alone is not the contract"}
- {id: logging.otel.success.exports_metric, module: logging, tier: P0, event: success, assertions: [exports_metric], exercised_on: [chat_completions, messages, responses, embeddings], source: "integrations/otel/logger.py", rationale: "OTEL spans on every call path"}
- {id: logging.otel.stream.exports_metric, module: logging, tier: P0, event: stream, assertions: [exports_metric], exercised_on: [chat_completions, messages, responses], source: "integrations/otel/logger.py", rationale: "Streaming closes the LLM span from the stream path; historically prone to duplicate/orphaned spans"}
- {id: logging.otel.stream.records_ttft, module: logging, tier: P1, event: stream, assertions: [records_ttft], exercised_on: [chat_completions, messages, responses], source: "integrations/otel/mappers/genai.py", rationale: "TTFT is the streaming latency SLI; a zero or span-length value silently corrupts dashboards"}

View file

@ -40,6 +40,7 @@ LlmEndpoint = Literal[
"audio_transcriptions",
"moderations",
"realtime",
"google_native",
]
LlmRoute = Literal[

View file

@ -124,7 +124,12 @@ class StreamingResponse(BaseModel):
non-streaming `application/json`), the response headers (lowercased names, e.g.
the x-ratelimit-* pacing headers and retry-after on a 429), 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."""
and streaming, where one validated JSON model does not fit.
`stream_done` records whether the terminal OpenAI `data: [DONE]` line arrived.
The consumed body is elided to "<streamed>", so that line is otherwise
unobservable, and routes differ on whether sending it is correct: OpenAI-shaped
streams must, while Google-native streams must not."""
status_code: int
call_id: str | None = None # x-litellm-call-id header
@ -134,6 +139,7 @@ class StreamingResponse(BaseModel):
body: str
chunks: int = 0 # streamed events (0 for non-streaming)
stream_events: list[str] = []
stream_done: bool = False
# First in-stream error event, if any. A streamed call commits its HTTP 200
# before the upstream completes, so upstream failures (e.g. insufficient
# quota) arrive as SSE error events inside an otherwise-successful response;

View file

@ -126,6 +126,19 @@ class ModerationRequest(BaseModel):
input: str
class GenerateContentPart(BaseModel):
text: str
class GenerateContentContent(BaseModel):
role: str = "user"
parts: list[GenerateContentPart]
class GenerateContentBody(BaseModel):
contents: list[GenerateContentContent]
class ResponsesOutputContent(BaseModel):
type: str | None = None
text: str | None = None
@ -423,6 +436,30 @@ class EndpointsClient:
response_type=ImagesResult,
)
def generate_content(self, key: str, model: str, text: str) -> StreamingResponse:
"""Google-native GenerateContent on a litellm-managed deployment.
This is the un-prefixed `/v1beta/models/...` route the google-genai and
Vertex SDKs talk to, not the `/gemini/...` or `/vertex_ai/...` passthrough.
"""
return self._send(
f"/v1beta/models/{model}:generateContent",
key,
GenerateContentBody(
contents=[GenerateContentContent(parts=[GenerateContentPart(text=text)])]
),
)
def stream_generate_content(self, key: str, model: str, text: str) -> StreamingResponse:
return self._send(
f"/v1beta/models/{model}:streamGenerateContent",
key,
GenerateContentBody(
contents=[GenerateContentContent(parts=[GenerateContentPart(text=text)])]
),
stream=True,
)
def build_endpoints_client(proxy: ProxyClient) -> EndpointsClient:
return EndpointsClient(proxy=proxy)

View file

@ -0,0 +1,133 @@
"""Live e2e: the Google-native GenerateContent surface on a managed deployment.
Customers point the google-genai and Vertex Java SDKs at the un-prefixed
`/v1beta/models/{model}:generateContent` route, so this route has to behave like
Google's own endpoint while still costing the call like every other litellm path.
Three shipped regressions live here, and each one is a separate customer symptom:
- the cost header was missing, so `generateContent` traffic could not be
reconciled against spend the way `/chat/completions` can (LIT-4076)
- streamed events were re-wrapped, producing a doubled `data: data:` prefix (and
at one point a literal Python `b'data:`), which no SSE client can parse
- the stream carried OpenAI's terminal `data: [DONE]` sentinel, which the Vertex
Java SDK rejects because Google never sends it
The streaming assertions only hold for a stream that actually succeeded: a
first-chunk upstream failure legitimately falls back to the OpenAI error shape
and does emit `[DONE]`, so the test proves real content arrived first.
"""
from __future__ import annotations
import pytest
from pydantic import BaseModel
from e2e_config import unique_marker
from e2e_http import StreamingResponse, require_successful_call
from endpoints_client import EndpointsClient
from lifecycle import ResourceManager
from models import LiteLLMParamsBody
pytestmark = pytest.mark.e2e
UPSTREAM_MODEL = "gemini/gemini-2.5-flash"
class _StreamPart(BaseModel):
text: str | None = None
class _StreamContent(BaseModel):
parts: list[_StreamPart] = []
class _StreamCandidate(BaseModel):
content: _StreamContent | None = None
class _StreamEvent(BaseModel):
candidates: list[_StreamCandidate] = []
def _managed_deployment(
client: EndpointsClient, resources: ResourceManager, label: str
) -> str:
model = f"e2e-google-native-{label}-{unique_marker()}"
model_id = client.create_model(
model,
LiteLLMParamsBody(model=UPSTREAM_MODEL, api_key="os.environ/GEMINI_API_KEY"),
)
resources.defer(lambda: client.delete_model(model_id))
return model
def _streamed_text(result: StreamingResponse) -> str:
"""Concatenate candidate text across events, validating each event parses.
A doubled `data:` prefix survives the harness parser as a payload that still
starts with `data:`, so validating every event is what turns that framing bug
into a failure here rather than a silently empty string.
"""
return "".join(
part.text
for event in result.stream_events
for candidate in _StreamEvent.model_validate_json(event).candidates
for part in (candidate.content.parts if candidate.content else [])
if part.text
)
class TestGoogleNativeGenerateContent:
@pytest.mark.covers("llm.google_native.gemini.basic.nonstream.cost_logged")
def test_generate_content_returns_response_cost_header(
self, endpoints_client: EndpointsClient, resources: ResourceManager
) -> None:
model = _managed_deployment(endpoints_client, resources, "cost")
key = resources.key()
result = endpoints_client.generate_content(
key, model, f"Reply with the single word ok. {unique_marker()}"
)
require_successful_call(result)
assert result.call_id, "generateContent must stamp x-litellm-call-id"
assert result.response_cost is not None, (
"generateContent returned no x-litellm-response-cost header; "
"google-native traffic cannot be reconciled against spend without it"
)
assert result.response_cost > 0, (
f"x-litellm-response-cost must be a real cost, got {result.response_cost}"
)
@pytest.mark.covers("llm.google_native.gemini.basic.stream.works")
def test_stream_generate_content_frames_sse_the_way_google_sdks_expect(
self, endpoints_client: EndpointsClient, resources: ResourceManager
) -> None:
model = _managed_deployment(endpoints_client, resources, "stream")
key = resources.key()
result = endpoints_client.stream_generate_content(
key, model, f"Count from one to five, one number per line. {unique_marker()}"
)
require_successful_call(result)
assert result.is_streaming, (
f"expected text/event-stream, got content-type {result.content_type!r}"
)
assert result.stream_error is None, f"stream carried an error: {result.stream_error}"
assert result.stream_events, f"stream delivered no data events (chunks={result.chunks})"
assert _streamed_text(result).strip(), "stream delivered events but no candidate text"
doubled = [event for event in result.stream_events if event.lstrip().startswith("data:")]
assert not doubled, (
f"{len(doubled)} event(s) carry a second data: prefix, so the proxy re-wrapped "
f"already-framed SSE; first offender: {doubled[0][:120]!r}"
)
leaked = [event for event in result.stream_events if event.startswith("b'")]
assert not leaked, (
f"event serialized as a Python bytes literal instead of text: {leaked[0][:120]!r}"
)
assert not result.stream_done, (
"google-native stream emitted the OpenAI [DONE] sentinel; Google never sends it "
"and the Vertex Java SDK rejects the stream when it appears"
)

View file

@ -0,0 +1,72 @@
"""Live e2e: the Prometheus request-queue-time histogram is actually emitted.
`litellm_request_queue_time_seconds` is how operators see how long a request
waited before the proxy started working on it, so it feeds saturation alerts. It
regressed to never being emitted at all (LIT-2034): the family was registered, so
a scrape still listed the metric, but no observation was ever recorded and every
dashboard built on it read empty.
That is why asserting the family exists is not enough. This drives a real call on
a uniquely-aliased key and then requires a sample carrying that alias with a
positive count, which is the part that stayed silent through the regression. The
histogram is written on the success-logging callback, so the scrape polls to a
deadline rather than sleeping once.
"""
from __future__ import annotations
import time
import pytest
from prometheus_client.parser import text_string_to_metric_families
from e2e_config import unique_marker
from lifecycle import ResourceManager
from logging_client import LoggingClient
pytestmark = pytest.mark.e2e
DRIVER_MODEL = "gemini-2.5-flash"
QUEUE_TIME_METRIC = "litellm_request_queue_time_seconds"
ALIAS_LABEL = "api_key_alias"
def _observation_count(exposition: str, alias: str) -> float | None:
"""The histogram's `_count` for our key's series, or None if never observed."""
for family in text_string_to_metric_families(exposition):
if family.name != QUEUE_TIME_METRIC:
continue
for sample in family.samples:
if sample.name == f"{QUEUE_TIME_METRIC}_count" and sample.labels.get(ALIAS_LABEL) == alias:
return sample.value
return None
class TestPrometheusRequestQueueTime:
@pytest.mark.covers("logging.prometheus.success.records_queue_time")
def test_queue_time_histogram_records_an_observation(
self, client: LoggingClient, resources: ResourceManager
) -> None:
alias = f"e2e-queue-time-{unique_marker()}"
key = client.key_with_alias(alias, models=[DRIVER_MODEL])
resources.defer(lambda: client.delete_key(key))
response = client.chat(key, DRIVER_MODEL, f"reply with one word {alias}")
assert response.model, f"driver call returned no model: {response}"
deadline = time.monotonic() + client.proxy.poll_timeout
count: float | None = None
while time.monotonic() < deadline:
count = _observation_count(client.scrape_metrics(), alias)
if count is not None and count > 0:
break
time.sleep(client.proxy.poll_interval)
assert count is not None, (
f"{QUEUE_TIME_METRIC} has no series for {ALIAS_LABEL}={alias}; the histogram was "
f"never observed for a request that succeeded"
)
assert count > 0, (
f"{QUEUE_TIME_METRIC} series for {alias} exists but recorded {count} observations; "
f"the metric is registered yet never written"
)