test(integration): streamed chat completions emit SSE keepalive pings while the upstream is silent before its first token (Pylon #7987)

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
kerry 2026-09-22 20:39:22 +00:00
parent c6c3881d7f
commit a1792250f8
2 changed files with 56 additions and 0 deletions

View file

@ -143,6 +143,9 @@
"tests/integration/streaming/test_stream_contracts.py::test_client_cancellation_releases_the_actual_provider_connection": [
"other.streaming.cancellation.closes_actual_provider_connection"
],
"tests/integration/streaming/test_ttft_keepalive.py::test_stream_emits_sse_ping_comments_before_the_first_data_frame_while_upstream_is_silent": [
"streaming.keepalive.sse_pings_fill_silent_time_to_first_token"
],
"tests/integration/routing/test_observed_routing.py::test_retry_counts_and_public_errors_match_actual_provider_attempts": [
"other.routing.retries.several_attempts_reach_success_without_hidden_retries",
"other.routing.errors.nonretryable_and_exhausted_failures_remain_errors"

View file

@ -0,0 +1,53 @@
import json
import time
import uuid
from collections.abc import Callable
from typing import Final
import pytest
from integration._support.client import Gateway
from integration._support.wire import Reply, Request, wire_server
from integration.streaming.test_stream_contracts import text_stream
UPSTREAM_SILENCE_SECONDS: Final = 3.0
KEEPALIVE_SECONDS: Final = 1
def _reply_after_silence(identity: str) -> Callable[[Request], Reply]:
def respond(_request: Request) -> Reply:
time.sleep(UPSTREAM_SILENCE_SECONDS)
return Reply(content_type="text/event-stream", chunks=text_stream(identity))
return respond
@pytest.mark.covers("streaming.keepalive.sse_pings_fill_silent_time_to_first_token")
def test_stream_emits_sse_ping_comments_before_the_first_data_frame_while_upstream_is_silent(
gateway: Gateway,
) -> None:
identity: Final = "stream-ttft-keepalive-" + uuid.uuid4().hex
with gateway.scenario() as scenario:
with wire_server(_reply_after_silence(identity)) as wire:
model: Final = scenario.model(api_base=wire.url + "/v1", keepalive_seconds=KEEPALIVE_SECONDS)
with gateway.client.stream(
"POST",
"/v1/chat/completions",
json={"model": model, "messages": [{"role": "user", "content": identity}], "stream": True},
headers={"Authorization": f"Bearer {gateway.key}"},
) as response:
assert response.status_code == 200, response.read().decode()
frames: Final = tuple(line for line in response.iter_lines() if line)
first_data: Final = next(index for index, line in enumerate(frames) if line.startswith("data:"))
assert first_data >= 1, f"No keepalive reached the client before the first data frame: {frames}"
assert frames[:first_data] == (": ping",) * first_data, frames
assert frames[-1] == "data: [DONE]", frames
deltas: Final = tuple(json.loads(line.removeprefix("data: ")) for line in frames[first_data:-1])
assert "".join(
choice["delta"].get("content", "") for chunk in deltas for choice in chunk["choices"]
) == "Hello 雪 café", frames
requests: Final = wire.drain()
assert len(requests) == 1
outbound: Final = json.loads(requests[0].body)
assert outbound["model"] == "gpt-4o-mini" and outbound["stream"] is True, outbound
assert "keepalive_seconds" not in outbound, outbound