test(e2e): open the least-busy stream under least-busy so its process counts it, and prove the busy deployment's health by draining to the terminator

This commit is contained in:
mateo-berri 2026-09-04 23:49:55 -07:00
parent 1bd70a6698
commit 308f66f114

View file

@ -19,22 +19,25 @@ then land on the fast one: any process meets the slow deployment at most once
before routing around it. The control call's timeout proves the slow deployment
was still routable, so the fast picks were latency's doing, not a cooldown's.
Least-busy reads live traffic and ignores weights, so its group carries the same
1/0 split: one long streaming request opened under simple-shuffle lands on the
weighted deployment (its head names it) and is held unread, and every short
least-busy call sent while it is in flight must land on one of three weight-0
deployments. Three of them rather than one because a proxy process counts
in-flight requests in its own memory, reads the shared count from Redis only on
its first look at a group, and releases a call's count in a success callback
that runs some time after the response leaves it, so a process can still count
the previous call or two against whichever deployment took them. With three
calls and three idle deployments, every process's view keeps some idle
deployment at zero, strictly below the one holding the stream, so no call can
tie with it and lose the tie on insertion order. The group gets no warm-up call
for the same reason: a process that served it before the stream opened would
route on its own stale copy, in which nothing is busy. The closing simple-shuffle
control call landing on the weighted deployment proves it was healthy the whole
time.
Least-busy reads live traffic, so its group of four equal deployments gets one
long streaming request, opened under least-busy and held unread (its head names
the deployment it landed on), and every short least-busy call sent while it is
in flight must land on one of the other three. The stream itself goes through
least-busy because a proxy process only starts counting in-flight requests once
it has routed a least-busy request, which is what registers the counting
callback, so a stream opened under another strategy would go uncounted in a
process that has never routed one. Three idle deployments rather than one
because a process counts in its own memory, reads the shared count from Redis
only on its first look at a group, and releases a call's count in a success
callback that runs some time after the response leaves it, so a process can
still count the previous call or two against whichever deployment took them;
with three calls and three idle deployments, every process's view keeps some
idle deployment at zero, strictly below the one holding the stream, so no call
can tie with it and lose the tie on insertion order. The group gets no warm-up
call for the same reason: a process that served it before the stream opened
would route on its own stale copy, in which nothing is busy. Draining the stream
to its terminator afterwards proves the deployment holding it was healthy the
whole time.
The per-request strategy comes in through `router_settings_override`, the same
knob a key or team's `router_settings` feeds, so one long-lived proxy configured
@ -46,7 +49,7 @@ from __future__ import annotations
import pytest
from complexity_router_client import ComplexityRouterClient
from e2e_config import unique_marker
from e2e_http import StreamHead
from e2e_http import StreamChunk, StreamHead, StreamStep, StreamTruncation
from lifecycle import ResourceManager
from models import LiteLLMParamsBody, ModelInfoBody, ModelNewBody, RouterSettingsOverride, RoutingStrategy
from reliability_support import REAL_KEY, REAL_MODEL, chat_override, model_id_of, open_chat_stream
@ -131,6 +134,15 @@ def _latency_picks(
return _latency_picks(client, key, group, slow, fast, (*history, _latency_pick(client, key, group, slow, fast)))
def _assert_streamed_to_the_end(drained: tuple[StreamStep, ...], busy: str | None) -> None:
truncations = [step for step in drained if isinstance(step, StreamTruncation)]
body = b"".join(step.data for step in drained if isinstance(step, StreamChunk))
assert not truncations and b"[DONE]" in body, (
f"the long stream on {busy} did not run to its terminator, so that deployment may not have been healthy: "
f"{truncations or body[-200:]!r}"
)
def _assert_shuffle_control_lands_on(client: ComplexityRouterClient, key: str, group: str, weighted: str) -> None:
control = _pick(client, key, group, "simple-shuffle")
assert control == weighted, (
@ -211,28 +223,27 @@ class TestReliabilityRoutingStrategies:
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
group = f"reliability-leastbusy-{unique_marker()}"
busy = _register(client, resources, group, _real(weight=1))
idle = frozenset(_register(client, resources, group, _real(weight=0)) for _ in range(STRATEGY_CALLS))
deployments = frozenset(_register(client, resources, group, _real(weight=1)) for _ in range(STRATEGY_CALLS + 1))
head = open_chat_stream(
client.proxy,
scoped_key,
group,
f"Write a 1500 word essay on the history of the telegraph. {unique_marker()}",
override=RouterSettingsOverride(routing_strategy="simple-shuffle"),
override=RouterSettingsOverride(routing_strategy="least-busy"),
max_tokens=3000,
)
assert isinstance(head, StreamHead), f"opening the long stream failed: {head}"
busy = head.headers.get("x-litellm-model-id")
try:
assert head.status_code == 200, f"the long stream should have opened with a 200, got {head.status_code}"
landed = head.headers.get("x-litellm-model-id")
assert landed == busy, f"the long stream landed on {landed!r}, not the weighted deployment {busy}"
assert busy in deployments, f"the long stream landed on {busy!r}, not one of {sorted(deployments)}"
idle = deployments - {busy}
picks = [_pick(client, scoped_key, group, "least-busy") for _ in range(STRATEGY_CALLS)]
assert all(pick in idle for pick in picks), (
f"least-busy picked {picks}, expected every call on one of {sorted(idle)} while {busy} still has the "
"long stream in flight"
)
_assert_shuffle_control_lands_on(client, scoped_key, group, busy)
finally:
for _ in head.steps:
pass
drained = tuple(head.steps)
_assert_streamed_to_the_end(drained, busy)