From 3da1a160609ee3e8c82f443722e8cd1f81659ed2 Mon Sep 17 00:00:00 2001 From: yucheng Date: Sat, 26 Sep 2026 00:43:37 +0000 Subject: [PATCH] test(integration): let gated replies wait as long as the spend-row polls Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- tests/integration/_support/wire.py | 5 ++++- .../test_passthrough_upstream_error_chaos.py | 6 +++++- .../test_passthrough_upstream_error_visibility.py | 8 +++++++- 3 files changed, 16 insertions(+), 3 deletions(-) diff --git a/tests/integration/_support/wire.py b/tests/integration/_support/wire.py index ed96d4e4e83..cb7c31ff836 100644 --- a/tests/integration/_support/wire.py +++ b/tests/integration/_support/wire.py @@ -28,6 +28,7 @@ class Reply: chunks: tuple[bytes, ...] | None = None abort_after: int | None = None gate_after_first: threading.Event | None = None + gate_timeout_seconds: float = 5 pause_between_chunks: float = 0 headers: Mapping[str, str] = MappingProxyType({}) @@ -88,7 +89,9 @@ def wire_server( self.wfile.write(b"%x\r\n%s\r\n" % (len(chunk), chunk)) self.wfile.flush() if index == 0 and reply.gate_after_first is not None: - assert reply.gate_after_first.wait(timeout=5), "Stream barrier was never released" + assert reply.gate_after_first.wait(timeout=reply.gate_timeout_seconds), ( + "Stream barrier was never released" + ) if reply.pause_between_chunks and index + 1 < len(reply.chunks): time.sleep(reply.pause_between_chunks) else: diff --git a/tests/integration/observability/test_passthrough_upstream_error_chaos.py b/tests/integration/observability/test_passthrough_upstream_error_chaos.py index 104d1191aae..a42eb576699 100644 --- a/tests/integration/observability/test_passthrough_upstream_error_chaos.py +++ b/tests/integration/observability/test_passthrough_upstream_error_chaos.py @@ -163,7 +163,11 @@ async def test_passthrough_disconnect_burst_logs_every_failure_once(gateway: Gat def respond(request: Request) -> Reply: if "streamGenerateContent" in request.target: return Reply( - status=429, content_type="text/event-stream", chunks=_RATE_LIMITED_FRAMES, gate_after_first=gate + status=429, + content_type="text/event-stream", + chunks=_RATE_LIMITED_FRAMES, + gate_after_first=gate, + gate_timeout_seconds=300, ) return Reply(status=200, body=json.dumps({"ok": True}).encode()) diff --git a/tests/integration/observability/test_passthrough_upstream_error_visibility.py b/tests/integration/observability/test_passthrough_upstream_error_visibility.py index d304cffa424..40c50eb4863 100644 --- a/tests/integration/observability/test_passthrough_upstream_error_visibility.py +++ b/tests/integration/observability/test_passthrough_upstream_error_visibility.py @@ -697,7 +697,13 @@ def test_gemini_passthrough_streaming_429_client_disconnect_still_logs_failure( frames: Final = (b'data: {"error":"rate limited"}\n\n', b"data: [DONE]\n\n") def respond(request: Request) -> Reply: - return Reply(status=429, content_type="text/event-stream", chunks=frames, gate_after_first=gate) + return Reply( + status=429, + content_type="text/event-stream", + chunks=frames, + gate_after_first=gate, + gate_timeout_seconds=120, + ) path: Final = tmp_path / "gemini-stream-429-disconnect.yaml" with wire_server(respond) as wire: