test(spend): build outage bursts with immutable generator sends

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
kerry 2026-10-03 01:51:28 +00:00
parent 1f701b92df
commit 79663c9037

View file

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