test(spend): make sink outage assertions race-free with a lone recovery probe

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
kerry 2026-10-03 00:07:48 +00:00
parent e4ea235bb3
commit 33d59fddae

View file

@ -205,44 +205,57 @@ def test_sink_outage_does_not_lose_spend_log_tags(gateway: Gateway, tmp_path: Pa
return {response.headers["x-litellm-call-id"] for response in responses}
first_ids: Final = call_ids(first)
first_landed: Final = _landed_tags(key, lambda values: len(values) == len(first))
assert len(first_landed) == len(first)
def events_for(ids: Set[str]) -> Set[str]:
def events_for(ids: Set[str]) -> Sequence[Mapping]:
events: Final = chain.from_iterable(json.loads(batch.body) for batch in delivered)
return {event["litellm_call_id"] for event in events if event.get("litellm_call_id") in ids}
return [event for event in events if event.get("litellm_call_id") in ids]
eventually(lambda: events_for(first_ids), lambda found: found == first_ids, seconds=70)
eventually(lambda: events_for(first_ids), lambda found: len(found) >= 1, seconds=70)
first_events: Final = events_for(first_ids)
first_occurrences: Final = [event["litellm_call_id"] for event in first_events]
assert set(first_occurrences) <= first_ids
assert len(first_occurrences) == len(set(first_occurrences)), "duplicate burst-1 delivery"
for event in first_events:
assert event["request_tags"] == EXPECTED
down.set()
with ThreadPoolExecutor(max_workers=5) as pool:
second: Final = tuple(chain.from_iterable(pool.map(lambda i: burst(100 + i), range(3))))
second_ids: Final = call_ids(second)
outage_probe: Final = eventually(
lambda: (len(rejected), events_for(second_ids)),
lambda: (len(rejected), {event["litellm_call_id"] for event in events_for(second_ids)}),
lambda state: state[0] >= 1 and state[1] == set(),
seconds=30,
)
assert outage_probe[0] >= 1, "sink saw no rejection during the outage window"
assert outage_probe[1] == set(), "burst-2 event delivered to a down sink"
down.clear()
third: Final = _tagged_requests(candidate, key, anthropic_model, openai_model, False, 200)
responses: Final = [*first, *second, *third]
assert all(response.status_code == 200 for response in responses), [
(response.status_code, response.text[:200]) for response in responses
]
third_ids: Final = call_ids(list(third))
recovery_probe: Final = next(iter(third_ids))
eventually(
lambda: events_for(third_ids),
lambda found: recovery_probe in found,
probe: Final = candidate.request(
"POST",
"/v1/chat/completions",
{
"model": openai_model,
"messages": [{"role": "user", "content": f"recovery probe {uuid.uuid4().hex}"}],
},
key=key,
headers=HEADERS,
)
assert probe.status_code == 200, probe.text
probe_id: Final = probe.headers["x-litellm-call-id"]
probe_events: Final = eventually(
lambda: events_for({probe_id}),
lambda found: len(found) >= 1,
seconds=70,
)
second_events: Final = chain.from_iterable(json.loads(batch.body) for batch in delivered)
second_occurrences: Final = [
event["litellm_call_id"] for event in second_events if event.get("litellm_call_id") in second_ids
]
assert len(probe_events) == 1, "recovery probe delivered to the sink more than once"
assert probe_events[0]["request_tags"] == EXPECTED
responses: Final = [*first, *second, probe]
second_events: Final = events_for(second_ids)
second_occurrences: Final = [event["litellm_call_id"] for event in second_events]
assert len(second_occurrences) == len(set(second_occurrences)), (
"duplicate burst-2 delivery after the outage"
)
assert events_for(second_ids) <= second_ids
_landed_tags(key, lambda values: len(values) == len(responses))