Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions tests/e2e/e2e_http.py
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,7 @@ class StreamingResponse(BaseModel):
headers: dict[str, str] = {}
body: str
chunks: int = 0 # streamed events (0 for non-streaming)
stream_events: list[str] = []
# 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;
Expand Down Expand Up @@ -296,10 +297,16 @@ def _streaming_outcome(resp: requests.Response, stream: bool) -> StreamingRespon
lines = cast("Iterator[bytes]", resp.iter_lines())
chunks = 0
stream_error: str | None = None
stream_events: list[str] = []
for line in lines:
if not line:
continue
chunks += 1
decoded_line = line.decode(errors="replace")
if decoded_line.startswith("data: "):
payload = decoded_line.removeprefix("data: ")
if payload != "[DONE]":
stream_events.append(payload)
if stream_error is None and (
line.startswith(b"event: error")
or b'"type":"error"' in line
Expand All @@ -315,6 +322,7 @@ def _streaming_outcome(resp: requests.Response, stream: bool) -> StreamingRespon
headers=headers,
body="<streamed>",
chunks=chunks,
stream_events=stream_events,
stream_error=stream_error,
)

Expand Down
34 changes: 30 additions & 4 deletions tests/e2e/llm_translation/endpoints_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
from __future__ import annotations

from dataclasses import dataclass
from typing import Literal

from pydantic import BaseModel

Expand All @@ -22,6 +23,7 @@ class ResponsesRequest(BaseModel):
model: str
input: str
instructions: str | None = None
stream: bool = False


class MessagesRequest(BaseModel):
Expand Down Expand Up @@ -100,6 +102,19 @@ def text(self) -> str:
)


class ResponsesStreamEvent(BaseModel):
event_id: str | None = None


class ResponsesStreamEventType(BaseModel):
type: str


class ResponsesOutputTextDeltaEvent(ResponsesStreamEvent):
type: Literal["response.output_text.delta"]
delta: str


class AnthropicContentBlock(BaseModel):
type: str | None = None
text: str | None = None
Expand Down Expand Up @@ -164,18 +179,29 @@ def create_model(self, model_name: str, litellm_params: LiteLLMParamsBody) -> st
def delete_model(self, model_id: str) -> None:
self.proxy.delete_model(model_id)

def _send(self, path: str, key: str, body: BaseModel) -> StreamingResponse:
def _send(
self, path: str, key: str, body: BaseModel, *, stream: bool = False
) -> StreamingResponse:
return self.proxy.transport.send(
path, headers=self.proxy.transport.bearer(key), json=body
path,
headers=self.proxy.transport.bearer(key),
json=body,
stream=stream,
)

def responses(self, key: str, model: str, text: str) -> StreamingResponse:
def responses(
self, key: str, model: str, text: str, *, stream: bool = False
) -> StreamingResponse:
return self._send(
"/v1/responses",
key,
ResponsesRequest(
model=model, input=text, instructions="You are a helpful assistant"
model=model,
input=text,
instructions="You are a helpful assistant",
stream=stream,
),
stream=stream,
)

def messages(
Expand Down
45 changes: 44 additions & 1 deletion tests/e2e/llm_translation/test_responses_e2e.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,17 +8,24 @@
from __future__ import annotations

import pytest
from pydantic import ValidationError

from e2e_config import unique_marker
from e2e_http import require_successful_call
from endpoints_client import EndpointsClient, ResponsesResult
from endpoints_client import (
EndpointsClient,
ResponsesOutputTextDeltaEvent,
ResponsesResult,
ResponsesStreamEventType,
)
from lifecycle import ResourceManager
from models import LiteLLMParamsBody

pytestmark = pytest.mark.e2e


class TestResponses:
@pytest.mark.covers("llm.responses.openai.basic.nonstream.works")
def test_responses_returns_completion(
self, endpoints_client: EndpointsClient, resources: ResourceManager
) -> None:
Expand All @@ -34,3 +41,39 @@ def test_responses_returns_completion(
require_successful_call(result)
parsed = ResponsesResult.model_validate_json(result.body)
assert parsed.text.strip(), f"/responses returned no output text: {result.body[:300]}"

@pytest.mark.covers("llm.responses.openai.basic.stream.works")
def test_responses_streaming_returns_completion(
self, endpoints_client: EndpointsClient, resources: ResourceManager
) -> None:
model = f"e2e-responses-{unique_marker()}"
model_id = endpoints_client.create_model(
model,
LiteLLMParamsBody(model="openai/gpt-4o-mini", api_key="os.environ/OPENAI_API_KEY"),
)
resources.defer(lambda: endpoints_client.delete_model(model_id))
key = resources.key()

result = endpoints_client.responses(key, model, "reply with one word", stream=True)
require_successful_call(result)
delta_events = tuple(
parsed
for event in result.stream_events
if (parsed := _parse_stream_event(event)) is not None
)

assert any(event.delta for event in delta_events), "responses stream returned no text deltas"
assert result.stream_events, "responses stream returned no events"
assert (
ResponsesStreamEventType.model_validate_json(result.stream_events[-1]).type
== "response.completed"
), "responses stream did not terminate with response.completed"


def _parse_stream_event(
event: str,
) -> ResponsesOutputTextDeltaEvent | None:
try:
return ResponsesOutputTextDeltaEvent.model_validate_json(event)
except ValidationError:
return None
Loading