diff --git a/tests/e2e/coverage_registry/logging.yaml b/tests/e2e/coverage_registry/logging.yaml index c4b26387c13..afb6dbc964e 100644 --- a/tests/e2e/coverage_registry/logging.yaml +++ b/tests/e2e/coverage_registry/logging.yaml @@ -9,6 +9,7 @@ - {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.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.failure.exports_metric, module: logging, tier: P0, event: failure, assertions: [exports_metric], exercised_on: [chat_completions, messages], source: "integrations/otel/logger.py", rationale: "Error spans for observability continuity"} - {id: logging.braintrust.success.logs_spend, module: logging, tier: P1, event: success, assertions: [logs_spend], exercised_on: [chat_completions, messages], source: "integrations/braintrust_logging.py", rationale: "Evals platform spend"} - {id: logging.langsmith.success.logs_spend, module: logging, tier: P1, event: success, assertions: [logs_spend], exercised_on: [chat_completions, messages], source: "integrations/langsmith.py", rationale: "LangChain ecosystem"} diff --git a/tests/e2e/e2e_http.py b/tests/e2e/e2e_http.py index ad53b2b4aa2..ff296969079 100644 --- a/tests/e2e/e2e_http.py +++ b/tests/e2e/e2e_http.py @@ -121,6 +121,11 @@ class StreamingResponse(BaseModel): headers: dict[str, str] = {} body: str chunks: int = 0 # streamed events (0 for non-streaming) + # 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; + # the consumed body is elided, so this is the only place they surface. + stream_error: str | None = None @property def ok(self) -> bool: @@ -289,7 +294,19 @@ def _streaming_outcome(resp: requests.Response, stream: bool) -> StreamingRespon body=resp.text, ) lines = cast("Iterator[bytes]", resp.iter_lines()) - chunks = sum(1 for line in lines if line) + chunks = 0 + stream_error: str | None = None + for line in lines: + if not line: + continue + chunks += 1 + if stream_error is None and ( + line.startswith(b"event: error") + or b'"type":"error"' in line + or b'"type": "error"' in line + or line.startswith(b'data: {"error"') + ): + stream_error = line.decode(errors="replace")[:300] return StreamingResponse( status_code=resp.status_code, call_id=call_id, @@ -298,6 +315,7 @@ def _streaming_outcome(resp: requests.Response, stream: bool) -> StreamingRespon headers=headers, body="", chunks=chunks, + stream_error=stream_error, ) diff --git a/tests/e2e/logging/logging_client.py b/tests/e2e/logging/logging_client.py index bffdf71ed80..e37d6175705 100644 --- a/tests/e2e/logging/logging_client.py +++ b/tests/e2e/logging/logging_client.py @@ -77,11 +77,12 @@ WEATHER_TOOL = ChatTool( class ResponsesRequestBody(BaseModel): - """OpenAI Responses API /v1/responses request (non-streaming).""" + """OpenAI Responses API /v1/responses request.""" model: str input: str max_output_tokens: int + stream: bool | None = None class TeamCallbackBody(BaseModel): @@ -464,30 +465,43 @@ class LoggingClient: json=body, ) - def messages_raw(self, key: str, model: str, text: str, *, max_tokens: int = 16) -> StreamingResponse: - """Non-streaming POST /v1/messages (Anthropic-native body): raw outcome - judged by status/body/headers, for tests that need x-litellm-call-id.""" + def messages_raw( + self, key: str, model: str, text: str, *, max_tokens: int = 16, stream: bool = False + ) -> StreamingResponse: + """POST /v1/messages (Anthropic-native body): raw outcome judged by + status/body/headers, for tests that need x-litellm-call-id. With + ``stream=True`` the SSE body is consumed and its events counted.""" + body = AnthropicMessagesBody( + model=model, + max_tokens=max_tokens, + messages=[ChatMessage(role="user", content=text)], + stream=True if stream else None, + ) + if stream: + return self.gateway.transport.stream( + "/v1/messages", headers=self.gateway.transport.bearer(key), json=body + ) return self.gateway.transport.send( - "/v1/messages", - headers=self.gateway.transport.bearer(key), - json=AnthropicMessagesBody( - model=model, - max_tokens=max_tokens, - messages=[ChatMessage(role="user", content=text)], - ), + "/v1/messages", headers=self.gateway.transport.bearer(key), json=body ) def responses_raw( - self, key: str, model: str, text: str, *, max_output_tokens: int = 64 + self, key: str, model: str, text: str, *, max_output_tokens: int = 64, stream: bool = False ) -> StreamingResponse: - """Non-streaming POST /v1/responses (OpenAI Responses API): raw outcome - judged by status/body/headers, for tests that need x-litellm-call-id. + """POST /v1/responses (OpenAI Responses API): raw outcome judged by + status/body/headers, for tests that need x-litellm-call-id. max_output_tokens caps reasoning-model output cost; a capped response is - still a 200 and still exports the trace.""" + still a 200 and still exports the trace. With ``stream=True`` the SSE + body is consumed and its events counted.""" + body = ResponsesRequestBody( + model=model, input=text, max_output_tokens=max_output_tokens, stream=True if stream else None + ) + if stream: + return self.gateway.transport.stream( + "/v1/responses", headers=self.gateway.transport.bearer(key), json=body + ) return self.gateway.transport.send( - "/v1/responses", - headers=self.gateway.transport.bearer(key), - json=ResponsesRequestBody(model=model, input=text, max_output_tokens=max_output_tokens), + "/v1/responses", headers=self.gateway.transport.bearer(key), json=body ) def scrape_metrics(self) -> str: diff --git a/tests/e2e/logging/test_otel_trace_e2e.py b/tests/e2e/logging/test_otel_trace_e2e.py index 90887ea5510..64df851fd5c 100644 --- a/tests/e2e/logging/test_otel_trace_e2e.py +++ b/tests/e2e/logging/test_otel_trace_e2e.py @@ -27,7 +27,7 @@ from e2e_config import CHEAP_ANTHROPIC_MODEL, CHEAP_OPENAI_MODEL, unique_marker from e2e_http import NoBody, StreamingResponse, require_successful_call from lifecycle import ResourceManager from logging_client import LoggingClient -from otel_client import JaegerTrace, OtelReader +from otel_client import JaegerSpan, JaegerTrace, OtelReader pytestmark = pytest.mark.e2e @@ -96,7 +96,9 @@ def _chain_reaches(span_id: str, root_id: str, trace: JaegerTrace) -> bool: return False -def _assert_complete_trace(hits: list[JaegerTrace], *, route: str, genai_span: str) -> None: +def _assert_complete_trace( + hits: list[JaegerTrace], *, route: str, genai_span: str, require_cost_span: bool = True +) -> None: """The enforced behavior: the destination holds exactly one trace for the call, rooted at the SERVER span, with auth/db/cost children and the gen-AI span all connected into that one tree - no dangling parent references.""" @@ -135,7 +137,8 @@ def _assert_complete_trace(hits: list[JaegerTrace], *, route: str, genai_span: s assert any(name.startswith(DB_SPAN_PREFIX) for name in names), ( f"no db ('{DB_SPAN_PREFIX}*') span in the trace; spans: {names}" ) - assert COST_SPAN in names, f"cost write span {COST_SPAN!r} missing; spans: {names}" + if require_cost_span: + assert COST_SPAN in names, f"cost write span {COST_SPAN!r} missing; spans: {names}" genai = next((span for span in trace.spans if span.operation_name == genai_span), None) assert genai is not None, f"gen-AI span {genai_span!r} missing; spans: {names}" @@ -146,8 +149,16 @@ def _assert_complete_trace(hits: list[JaegerTrace], *, route: str, genai_span: s ) -def _settled_names(*, route: str, genai_span: str) -> set[str]: - return {f"POST {route}", f"auth {route}", COST_SPAN, genai_span} +def _settled_names(*, route: str, genai_span: str, require_cost_span: bool = True) -> set[str]: + names = {f"POST {route}", f"auth {route}", genai_span} + return (names | {COST_SPAN}) if require_cost_span else names + + +def _tag(span: JaegerSpan, key: str) -> str | int | float | bool | None: + for tag in span.tags: + if tag.key == key: + return tag.value + return None class TestOtelTraceCompleteness: @@ -257,3 +268,179 @@ class TestOtelTraceCompleteness: settled_prefixes={DB_SPAN_PREFIX}, ) _assert_complete_trace(hits, route=route, genai_span=genai_span) + + @pytest.mark.covers("logging.otel.stream.exports_metric", exercised_on=["chat_completions"]) + def test_chat_completions_stream_exports_complete_trace( + self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager + ) -> None: + """A successful streamed `/chat/completions` request should export one + complete OTEL trace. The trace must contain a single root `SERVER` + span, with the auth, database, cost, and gen-AI `CLIENT` spans all + connected back to that root. + + Streaming has an additional lifecycle risk because the gen-AI span is + closed by the stream-consumption path after the final chunk has + arrived and usage has been aggregated. Historically, this has caused + duplicate or orphaned spans. + + The test therefore confirms that: + + * The response actually streams. + * Exactly one gen-AI span is created for the request. + * The gen-AI span contains `litellm.request.streaming=true`. + """ + route = "/chat/completions" + _assert_otel_destination_configured(client) + + key = client.key_with_alias(f"otel-stream-chat-{unique_marker()}", models=[MODEL]) + resources.defer(lambda: client.delete_key(key)) + + marker = unique_marker() + outcome = _first_ok( + client, + lambda: client.chat_raw(key, MODEL, f"reply with one word {marker}", stream=True, max_tokens=16), + ) + assert outcome.call_id is not None, "success response must carry x-litellm-call-id" + 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" + assert outcome.stream_error is None, ( + f"the stream carried an upstream error event despite the 200: {outcome.stream_error}" + ) + + genai_span = f"chat {MODEL}" + hits = otel_reader.poll_traces_for_call( + call_id=outcome.call_id, + settled_names=_settled_names(route=route, genai_span=genai_span), + settled_prefixes={DB_SPAN_PREFIX}, + ) + _assert_complete_trace(hits, route=route, genai_span=genai_span) + + genai_spans = [span for span in hits[0].spans if span.operation_name == genai_span] + assert len(genai_spans) == 1, ( + f"a streamed call must produce exactly ONE gen-AI span, got {len(genai_spans)}; " + f"spans: {hits[0].span_names()}" + ) + assert _tag(genai_spans[0], "litellm.request.streaming") is True, ( + "the gen-AI span must record litellm.request.streaming=true; its absence means " + "the stream flag was dropped before the model call" + ) + + @pytest.mark.covers("logging.otel.stream.exports_metric", exercised_on=["messages"]) + def test_messages_stream_exports_complete_trace( + self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager + ) -> None: + """A successful streamed `/v1/messages` request should export one + complete OTEL trace. The trace must contain a single root `SERVER` + span, with the auth, database, cost, and gen-AI `CLIENT` spans all + connected back to that root. + + This endpoint has the same streaming lifecycle risk as + `/chat/completions`: the gen-AI span is closed by the + stream-consumption path after the final chunk has arrived and usage + has been aggregated. + + The test therefore confirms that: + + * The response actually streams. + * Exactly one gen-AI span is created for the request. + * The gen-AI span contains `litellm.request.streaming=true`. + """ + route = "/v1/messages" + _assert_otel_destination_configured(client) + + key = client.key_with_alias(f"otel-stream-messages-{unique_marker()}", models=[MODEL]) + resources.defer(lambda: client.delete_key(key)) + + marker = unique_marker() + outcome = _first_ok( + client, + lambda: client.messages_raw(key, MODEL, f"reply with one word {marker}", max_tokens=16, stream=True), + ) + assert outcome.call_id is not None, "success response must carry x-litellm-call-id" + 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" + assert outcome.stream_error is None, ( + f"the stream carried an upstream error event despite the 200: {outcome.stream_error}" + ) + + genai_span = f"chat {MODEL}" + hits = otel_reader.poll_traces_for_call( + call_id=outcome.call_id, + settled_names=_settled_names(route=route, genai_span=genai_span), + settled_prefixes={DB_SPAN_PREFIX}, + ) + _assert_complete_trace(hits, route=route, genai_span=genai_span) + + genai_spans = [span for span in hits[0].spans if span.operation_name == genai_span] + assert len(genai_spans) == 1, ( + f"a streamed call must produce exactly ONE gen-AI span, got {len(genai_spans)}; " + f"spans: {hits[0].span_names()}" + ) + assert _tag(genai_spans[0], "litellm.request.streaming") is True, ( + "the gen-AI span must record litellm.request.streaming=true; its absence means " + "the stream flag was dropped before the model call" + ) + + @pytest.mark.covers("logging.otel.stream.exports_metric", exercised_on=["responses"]) + def test_responses_stream_exports_complete_trace( + self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager + ) -> None: + """A successful streamed /v1/responses request should export one complete + OTEL trace. The trace must contain a single root SERVER span, with the + auth, database, and gen-AI CLIENT spans all connected back to that + root. + + This endpoint has the same streaming lifecycle risk as the other + streaming surfaces: the gen-AI span is closed by the + stream-consumption path after the final event has arrived and usage + has been aggregated. + + The test therefore confirms that: + + * The response actually streams. + * Exactly one gen-AI span is created for the request. + * Spend is recorded correctly. + """ + route = "/v1/responses" + _assert_otel_destination_configured(client) + + key = client.key_with_alias( + f"otel-stream-responses-{unique_marker()}", models=[CHEAP_OPENAI_MODEL] + ) + resources.defer(lambda: client.delete_key(key)) + + marker = unique_marker() + outcome = _first_ok( + client, + lambda: client.responses_raw(key, CHEAP_OPENAI_MODEL, f"reply with one word {marker}", stream=True), + ) + assert outcome.call_id is not None, "success response must carry x-litellm-call-id" + 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" + assert outcome.stream_error is None, ( + f"the stream carried an upstream error event despite the 200: {outcome.stream_error}" + ) + + genai_span = f"chat {CHEAP_OPENAI_MODEL}" + hits = otel_reader.poll_traces_for_call( + call_id=outcome.call_id, + settled_names=_settled_names(route=route, genai_span=genai_span, require_cost_span=False), + settled_prefixes={DB_SPAN_PREFIX}, + ) + _assert_complete_trace(hits, route=route, genai_span=genai_span, require_cost_span=False) + + genai_spans = [span for span in hits[0].spans if span.operation_name == genai_span] + assert len(genai_spans) == 1, ( + f"a streamed call must produce exactly ONE gen-AI span, got {len(genai_spans)}; " + f"spans: {hits[0].span_names()}" + ) + + spend_row = client.poll_proxy_spend_for_key(key) + assert spend_row is not None and spend_row.spend is not None and spend_row.spend > 0, ( + "a successful streamed responses call must record a positive-spend row in /spend/logs " + "(the cost-write SPAN is knowingly absent on this surface, LIT-4428, but the spend " + f"itself must land); got {spend_row!r}" + ) + assert spend_row.call_type == "aresponses", ( + f"the spend row must be attributed to the responses call type, got {spend_row.call_type!r}" + ) diff --git a/tests/e2e/models.py b/tests/e2e/models.py index bf90426188f..4140967f3e0 100644 --- a/tests/e2e/models.py +++ b/tests/e2e/models.py @@ -154,6 +154,7 @@ class AnthropicMessagesBody(BaseModel): model: str messages: list[ChatMessage] max_tokens: int + stream: bool | None = None class AnthropicMessagesResponse(BaseModel):