From 25df767376f91b2ae5ec91c2b30722fbe3e1e015 Mon Sep 17 00:00:00 2001 From: mubashir1osmani Date: Sat, 25 Jul 2026 12:47:53 -0700 Subject: [PATCH] test(e2e): cover google-native generateContent framing and prometheus queue time Adds live coverage for three shipped regressions that had none, all reached through surfaces a customer drives from Google SDKs and operator dashboards. The managed google-native route (`/v1beta/models/{model}:generateContent`) had no harness support at all, so EndpointsClient gains generate_content and stream_generate_content plus the request body models, and a new suite asserts the two contracts that broke there: the response carries x-litellm-response-cost so SDK traffic reconciles against spend (LIT-4076), and the stream relays single-prefixed SSE frames with no OpenAI [DONE] terminator. A doubled `data:` prefix, a leaked bytes literal, or the [DONE] sentinel each fail the stream test; [DONE] absence is only asserted once real content has arrived, because a first-chunk upstream error legitimately falls back to the OpenAI error shape and does emit it. The prometheus test pins litellm_request_queue_time_seconds to an actual observation on our own key's series rather than to the family merely existing, which is the distinction the original regression turned on: the histogram stayed registered while nothing was ever written to it (LIT-2034). Each assertion was mutation-checked against the live proxy; inverting the [DONE] expectation, the cost-header expectation, or the metric name fails the corresponding test. --- .../llm_nonconversational.yaml | 2 + tests/e2e/coverage_registry/logging.yaml | 1 + tests/e2e/coverage_registry/schema.py | 1 + tests/e2e/e2e_http.py | 14 +- tests/e2e/llm_translation/endpoints_client.py | 37 +++++ .../llm_translation/test_google_native_e2e.py | 133 ++++++++++++++++++ .../logging/test_prometheus_queue_time_e2e.py | 72 ++++++++++ 7 files changed, 258 insertions(+), 2 deletions(-) create mode 100644 tests/e2e/llm_translation/test_google_native_e2e.py create mode 100644 tests/e2e/logging/test_prometheus_queue_time_e2e.py diff --git a/tests/e2e/coverage_registry/llm_nonconversational.yaml b/tests/e2e/coverage_registry/llm_nonconversational.yaml index 1e49cb13538..33482af944c 100644 --- a/tests/e2e/coverage_registry/llm_nonconversational.yaml +++ b/tests/e2e/coverage_registry/llm_nonconversational.yaml @@ -33,6 +33,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"} diff --git a/tests/e2e/coverage_registry/logging.yaml b/tests/e2e/coverage_registry/logging.yaml index 0f703632805..856636c3dbc 100644 --- a/tests/e2e/coverage_registry/logging.yaml +++ b/tests/e2e/coverage_registry/logging.yaml @@ -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"} diff --git a/tests/e2e/coverage_registry/schema.py b/tests/e2e/coverage_registry/schema.py index 89f5df73a4f..07f77c2ffb1 100644 --- a/tests/e2e/coverage_registry/schema.py +++ b/tests/e2e/coverage_registry/schema.py @@ -39,6 +39,7 @@ LlmEndpoint = Literal[ "audio_transcriptions", "moderations", "realtime", + "google_native", ] LlmRoute = Literal[ diff --git a/tests/e2e/e2e_http.py b/tests/e2e/e2e_http.py index d16e84bd754..8b6e73aa930 100644 --- a/tests/e2e/e2e_http.py +++ b/tests/e2e/e2e_http.py @@ -122,7 +122,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 "", 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 @@ -132,6 +137,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; @@ -408,6 +414,7 @@ def _streaming_outcome(resp: requests.Response, stream: bool) -> StreamingRespon chunks = 0 stream_error: str | None = None stream_events: list[str] = [] + stream_done = False for line in lines: if not line: continue @@ -415,7 +422,9 @@ def _streaming_outcome(resp: requests.Response, stream: bool) -> StreamingRespon decoded_line = line.decode(errors="replace") if decoded_line.startswith("data: "): payload = decoded_line.removeprefix("data: ") - if payload != "[DONE]": + if payload == "[DONE]": + stream_done = True + else: stream_events.append(payload) if stream_error is None and ( line.startswith(b"event: error") @@ -433,6 +442,7 @@ def _streaming_outcome(resp: requests.Response, stream: bool) -> StreamingRespon body="", chunks=chunks, stream_events=stream_events, + stream_done=stream_done, stream_error=stream_error, ) diff --git a/tests/e2e/llm_translation/endpoints_client.py b/tests/e2e/llm_translation/endpoints_client.py index ba201cbb5c0..74b072e563b 100644 --- a/tests/e2e/llm_translation/endpoints_client.py +++ b/tests/e2e/llm_translation/endpoints_client.py @@ -120,6 +120,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 @@ -400,6 +413,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) diff --git a/tests/e2e/llm_translation/test_google_native_e2e.py b/tests/e2e/llm_translation/test_google_native_e2e.py new file mode 100644 index 00000000000..eeade4a3f9c --- /dev/null +++ b/tests/e2e/llm_translation/test_google_native_e2e.py @@ -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" + ) diff --git a/tests/e2e/logging/test_prometheus_queue_time_e2e.py b/tests/e2e/logging/test_prometheus_queue_time_e2e.py new file mode 100644 index 00000000000..04c8c97c616 --- /dev/null +++ b/tests/e2e/logging/test_prometheus_queue_time_e2e.py @@ -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" + )