From 79663c9037461c10d24fe8d30666d79f47db1eac Mon Sep 17 00:00:00 2001 From: kerry Date: Sat, 3 Oct 2026 01:51:28 +0000 Subject: [PATCH] test(spend): build outage bursts with immutable generator sends Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../test_spend_log_request_tags_chaos.py | 50 +++++++++---------- 1 file changed, 23 insertions(+), 27 deletions(-) diff --git a/tests/integration/spend/test_spend_log_request_tags_chaos.py b/tests/integration/spend/test_spend_log_request_tags_chaos.py index 80df8ac63cd..068e2917e1a 100644 --- a/tests/integration/spend/test_spend_log_request_tags_chaos.py +++ b/tests/integration/spend/test_spend_log_request_tags_chaos.py @@ -1,13 +1,14 @@ import json import threading import uuid -from collections.abc import Mapping, Sequence, Set +from collections.abc import Callable, Mapping, Sequence, Set from concurrent.futures import ThreadPoolExecutor from hashlib import sha256 from itertools import chain from pathlib import Path from typing import Final +import httpx import psutil import pytest from integration._support.client import Gateway, eventually @@ -205,17 +206,27 @@ def test_sink_outage_does_not_lose_spend_log_tags(gateway: Gateway, tmp_path: Pa events: Final = chain.from_iterable(json.loads(batch.body) for batch in batches) return [event for event in events if event.get("litellm_call_id") in ids] - first: Final = [] - for index in range(3): - first_sends: Final = _requests( - candidate, key, anthropic_model, openai_model, False, f"burst {index} {uuid.uuid4().hex}" + def send_and_await(send: Callable[[], httpx.Response], batches: Sequence[Request]) -> httpx.Response: + response: Final = send() + assert response.status_code == 200, response.text + call_id: Final = response.headers["x-litellm-call-id"] + eventually(lambda: events_for(batches, {call_id}), lambda found: len(found) >= 1, seconds=30) + return response + + def burst_sends(label: str): + return chain.from_iterable( + _requests( + candidate, + key, + anthropic_model, + openai_model, + False, + f"{label} {index} {uuid.uuid4().hex}", + ) + for index in range(3) ) - for send in first_sends: - response = send() - assert response.status_code == 200, response.text - first.append(response) - call_id = response.headers["x-litellm-call-id"] - eventually(lambda: events_for(delivered, {call_id}), lambda found: len(found) == 1, seconds=30) + + first: Final = tuple(send_and_await(send, delivered) for send in burst_sends("burst")) first_ids: Final = {response.headers["x-litellm-call-id"] for response in first} first_events: Final = events_for(delivered, first_ids) first_occurrences: Final = [event["litellm_call_id"] for event in first_events] @@ -223,22 +234,7 @@ def test_sink_outage_does_not_lose_spend_log_tags(gateway: Gateway, tmp_path: Pa for event in first_events: assert event["request_tags"] == EXPECTED down.set() - second: Final = [] - for index in range(3): - second_sends: Final = _requests( - candidate, - key, - anthropic_model, - openai_model, - False, - f"outage {index} {uuid.uuid4().hex}", - ) - for send in second_sends: - response = send() - assert response.status_code == 200, response.text - second.append(response) - call_id = response.headers["x-litellm-call-id"] - eventually(lambda: events_for(rejected, {call_id}), lambda found: len(found) >= 1, seconds=30) + second: Final = tuple(send_and_await(send, rejected) for send in burst_sends("outage")) second_ids: Final = {response.headers["x-litellm-call-id"] for response in second} rejected_ids: Final = {event["litellm_call_id"] for event in events_for(rejected, second_ids)} assert rejected_ids == second_ids, "outage burst was not rejected by the down sink"