test(integration): let gated replies wait as long as the spend-row polls
Some checks failed
LiteLLM Rust / rust-lint (push) Waiting to run
LiteLLM Rust / rust-test (push) Waiting to run
LiteLLM Rust / rust-wheel (push) Waiting to run
Terraform Provider / gofmt, vet, build, test (push) Has been cancelled
Terraform Provider / Provider endpoints vs proxy OpenAPI schema (push) Has been cancelled

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
yucheng 2026-09-26 00:43:37 +00:00
parent 6a4a17f013
commit 3da1a16060
3 changed files with 16 additions and 3 deletions

View file

@ -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:

View file

@ -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())

View file

@ -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: