mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-08 03:08:45 +00:00
test(e2e): tag the remaining quota_management tests and record budget and spend client steps (#44966)
* test(e2e): add enum values, auto-discovering label gates and secret hiding for e2e metadata * test(e2e): tag quota_management tests with Subject metadata and record budget client steps * docs(e2e): name every markerless harness test file that carries no Subject * test(e2e): keep the step discovery comprehensions to one for clause
This commit is contained in:
parent
c943d650f4
commit
3538e87e45
10 changed files with 184 additions and 3 deletions
|
|
@ -17,6 +17,7 @@ from datetime import datetime
|
|||
from pydantic import AliasPath, BaseModel, Field, RootModel
|
||||
|
||||
from e2e_http import NoBody, Result, StreamingResponse, Success, unwrap
|
||||
from e2e_metadata import step
|
||||
from proxy_client import ProxyClient
|
||||
from models import (
|
||||
AnthropicMessagesBody,
|
||||
|
|
@ -253,15 +254,18 @@ class BudgetClient:
|
|||
)
|
||||
)
|
||||
|
||||
@step("Delete the virtual key")
|
||||
def delete_key(self, key: str) -> None:
|
||||
self.proxy.delete_key(key)
|
||||
|
||||
@step("Read the key's budget windows from /key/info")
|
||||
def key_budget_windows(self, key: str) -> list[BudgetWindowState]:
|
||||
"""A key's budget_limits windows as /key/info stores them. Each window's
|
||||
reset_at is advanced by the reset job in the same pass that zeroes the
|
||||
window's spend counter, so a strictly-later value proves the wipe ran."""
|
||||
return self.proxy.key_info(key).budget_limits or []
|
||||
|
||||
@step("Read the team's budget windows from /team/info")
|
||||
def team_budget_windows(self, team_id: str) -> list[BudgetWindowState]:
|
||||
"""Team analog of key_budget_windows, read from /team/info."""
|
||||
match self._team_info(team_id):
|
||||
|
|
@ -270,11 +274,13 @@ class BudgetClient:
|
|||
case _:
|
||||
return []
|
||||
|
||||
@step("Delete the end users {user_ids}")
|
||||
def delete_customers(self, user_ids: list[str]) -> None:
|
||||
self.proxy.delete_customers(user_ids)
|
||||
|
||||
# ---- chat (raw HTTP outcome: a budget block surfaces as a non-2xx) --
|
||||
|
||||
@step('Send a /chat/completions request to {model} with the prompt "{content}"')
|
||||
def chat(
|
||||
self,
|
||||
key: str,
|
||||
|
|
@ -297,6 +303,7 @@ class BudgetClient:
|
|||
),
|
||||
)
|
||||
|
||||
@step('Send a /v1/messages request to {model} with the prompt "{content}"')
|
||||
def messages(
|
||||
self,
|
||||
key: str,
|
||||
|
|
@ -317,6 +324,7 @@ class BudgetClient:
|
|||
|
||||
# ---- internal user --------------------------------------------------
|
||||
|
||||
@step("Create an internal user with max budget: {max_budget}")
|
||||
def create_user(self, *, max_budget: float, budget_duration: str | None = None) -> str:
|
||||
return unwrap(
|
||||
self.proxy.transport.post(
|
||||
|
|
@ -327,6 +335,7 @@ class BudgetClient:
|
|||
)
|
||||
).user_id
|
||||
|
||||
@step("Delete the internal user")
|
||||
def delete_user(self, user_id: str) -> None:
|
||||
_ = self.proxy.transport.post(
|
||||
"/user/delete",
|
||||
|
|
@ -335,6 +344,7 @@ class BudgetClient:
|
|||
response_type=NoBody,
|
||||
)
|
||||
|
||||
@step("Read the internal user's spend and budget from /user/info")
|
||||
def user_info(self, user_id: str) -> UserInfoRow | None:
|
||||
result = self.proxy.transport.get(
|
||||
"/user/info",
|
||||
|
|
@ -350,6 +360,7 @@ class BudgetClient:
|
|||
|
||||
# ---- customer / end-user -------------------------------------------
|
||||
|
||||
@step("Create the end user {customer_id}")
|
||||
def create_customer(
|
||||
self,
|
||||
customer_id: str,
|
||||
|
|
@ -369,6 +380,7 @@ class BudgetClient:
|
|||
|
||||
# ---- organization ---------------------------------------------------
|
||||
|
||||
@step("Create the organization {alias} with max budget: {max_budget}")
|
||||
def create_org(self, *, max_budget: float, alias: str, budget_duration: str | None = None) -> str:
|
||||
return unwrap(
|
||||
self.proxy.transport.post(
|
||||
|
|
@ -383,6 +395,7 @@ class BudgetClient:
|
|||
)
|
||||
).organization_id
|
||||
|
||||
@step("Read the organization's budget id from /organization/info")
|
||||
def org_budget_id(self, org_id: str) -> str | None:
|
||||
"""The id of the budget row backing an org; its budget_reset_at is read via
|
||||
budget_info (LIT-4570: /organization/new stores budget_duration without
|
||||
|
|
@ -399,6 +412,7 @@ class BudgetClient:
|
|||
case _:
|
||||
return None
|
||||
|
||||
@step("Delete the organization")
|
||||
def delete_org(self, org_id: str) -> None:
|
||||
_ = self.proxy.transport.delete(
|
||||
"/organization/delete",
|
||||
|
|
@ -409,6 +423,7 @@ class BudgetClient:
|
|||
|
||||
# ---- team -----------------------------------------------------------
|
||||
|
||||
@step("Create the team {alias} and wait until /team/info returns it")
|
||||
def create_team(
|
||||
self,
|
||||
*,
|
||||
|
|
@ -435,6 +450,7 @@ class BudgetClient:
|
|||
self._wait_for_team(team_id)
|
||||
return team_id
|
||||
|
||||
@step("Delete the team")
|
||||
def delete_team(self, team_id: str) -> None:
|
||||
_ = self.proxy.transport.post(
|
||||
"/team/delete",
|
||||
|
|
@ -463,6 +479,7 @@ class BudgetClient:
|
|||
assert last is not None
|
||||
raise AssertionError(last)
|
||||
|
||||
@step("Add the internal user to the team")
|
||||
def add_team_member(self, team_id: str, user_id: str, *, max_budget_in_team: float | None = None) -> None:
|
||||
last_body = ""
|
||||
for attempt in range(_TEAM_READY_ATTEMPTS):
|
||||
|
|
@ -484,6 +501,7 @@ class BudgetClient:
|
|||
break
|
||||
raise AssertionError(last_body)
|
||||
|
||||
@step("Update the team member's budget with /team/member_update")
|
||||
def update_team_member(
|
||||
self,
|
||||
team_id: str,
|
||||
|
|
@ -504,6 +522,7 @@ class BudgetClient:
|
|||
)
|
||||
assert resp.ok, resp.body
|
||||
|
||||
@step("Read the team member's budget reset time from /team/info")
|
||||
def member_budget_reset_at(self, team_id: str, user_id: str) -> str | None:
|
||||
"""The member's per-team budget_reset_at as /team/info reports it, or None if
|
||||
no reset is scheduled. The reset job advances this each time the window
|
||||
|
|
@ -519,6 +538,7 @@ class BudgetClient:
|
|||
|
||||
# ---- tag ------------------------------------------------------------
|
||||
|
||||
@step("Create the tag {name} with max budget: {max_budget}")
|
||||
def create_tag(self, name: str, *, max_budget: float) -> str:
|
||||
resp = self.proxy.transport.send(
|
||||
"/tag/new",
|
||||
|
|
@ -528,6 +548,7 @@ class BudgetClient:
|
|||
assert resp.ok, resp.body
|
||||
return name
|
||||
|
||||
@step("Delete the tag {name}")
|
||||
def delete_tag(self, name: str) -> None:
|
||||
_ = self.proxy.transport.post(
|
||||
"/tag/delete",
|
||||
|
|
@ -538,6 +559,7 @@ class BudgetClient:
|
|||
|
||||
# ---- model access group ---------------------------------------------
|
||||
|
||||
@step("Set a shared budget on the model access group {access_group}")
|
||||
def set_access_group_budget(
|
||||
self,
|
||||
access_group: str,
|
||||
|
|
@ -561,6 +583,7 @@ class BudgetClient:
|
|||
)
|
||||
)
|
||||
|
||||
@step("Read the budget and spend of the model access group {access_group}")
|
||||
def access_group_budget(self, access_group: str) -> AccessGroupBudgetResponse:
|
||||
return unwrap(
|
||||
self.proxy.transport.get(
|
||||
|
|
@ -571,6 +594,7 @@ class BudgetClient:
|
|||
)
|
||||
)
|
||||
|
||||
@step("Delete the budget on the model access group {access_group}")
|
||||
def delete_access_group_budget(self, access_group: str) -> None:
|
||||
_ = self.proxy.transport.delete(
|
||||
f"/access_group/{access_group}/budget",
|
||||
|
|
@ -581,6 +605,7 @@ class BudgetClient:
|
|||
|
||||
# ---- budget table ---------------------------------------------------
|
||||
|
||||
@step("Create a budget with /budget/new")
|
||||
def create_budget(
|
||||
self,
|
||||
*,
|
||||
|
|
@ -603,6 +628,7 @@ class BudgetClient:
|
|||
)
|
||||
).budget_id
|
||||
|
||||
@step("Delete the budget")
|
||||
def delete_budget(self, budget_id: str) -> None:
|
||||
_ = self.proxy.transport.post(
|
||||
"/budget/delete",
|
||||
|
|
@ -611,6 +637,7 @@ class BudgetClient:
|
|||
response_type=NoBody,
|
||||
)
|
||||
|
||||
@step("Read the budget from /budget/info")
|
||||
def budget_info(self, budget_id: str) -> tuple[BudgetRow, ...]:
|
||||
result = self.proxy.transport.post(
|
||||
"/budget/info",
|
||||
|
|
|
|||
|
|
@ -21,6 +21,7 @@ import time
|
|||
import pytest
|
||||
from e2e_config import unique_marker
|
||||
from e2e_http import StreamingResponse, require_successful_call
|
||||
from e2e_metadata import Domain, Mode, Provider, Subject, meta
|
||||
from quota_client import QuotaClient
|
||||
|
||||
pytestmark = pytest.mark.e2e
|
||||
|
|
@ -75,11 +76,27 @@ def _assert_blocked_inside_window(
|
|||
|
||||
class TestModelGroupAliasRateLimit:
|
||||
@pytest.mark.covers("quota_management.ratelimit.model_group_alias.shares_bucket")
|
||||
@meta(
|
||||
Subject(
|
||||
domain=Domain.SPEND_BUDGETS,
|
||||
providers=(Provider.ANTHROPIC,),
|
||||
models=(MODEL_GROUP, MODEL_ALIAS),
|
||||
mode=Mode.NONSTREAM,
|
||||
)
|
||||
)
|
||||
def test_alias_shares_rpm_bucket_with_model_group(self, client: QuotaClient, scoped_key: str) -> None:
|
||||
opened_at = _exhaust_rpm(client, scoped_key, MODEL_GROUP)
|
||||
_assert_blocked_inside_window(client, scoped_key, MODEL_ALIAS, opened_at)
|
||||
|
||||
@pytest.mark.covers("quota_management.ratelimit.model_group_alias.shares_bucket")
|
||||
@meta(
|
||||
Subject(
|
||||
domain=Domain.SPEND_BUDGETS,
|
||||
providers=(Provider.ANTHROPIC,),
|
||||
models=(MODEL_GROUP, MODEL_ALIAS),
|
||||
mode=Mode.NONSTREAM,
|
||||
)
|
||||
)
|
||||
def test_model_group_shares_rpm_bucket_with_alias(self, client: QuotaClient, scoped_key: str) -> None:
|
||||
opened_at = _exhaust_rpm(client, scoped_key, MODEL_ALIAS)
|
||||
_assert_blocked_inside_window(client, scoped_key, MODEL_GROUP, opened_at)
|
||||
|
|
|
|||
|
|
@ -36,6 +36,7 @@ from pydantic import BaseModel, RootModel
|
|||
|
||||
from e2e_config import unique_marker
|
||||
from e2e_http import Success
|
||||
from e2e_metadata import step
|
||||
from lifecycle import ResourceManager
|
||||
from models import LiteLLMParamsBody, SpendLogsParams
|
||||
from proxy_client import ProxyClient
|
||||
|
|
@ -133,6 +134,7 @@ def assert_fresh_tokens_billed_at(row: CostRow, input_rate: float) -> None:
|
|||
)
|
||||
|
||||
|
||||
@step("Wait for the request's cost breakdown in /spend/logs")
|
||||
def poll_cost_row(proxy: ProxyClient, request_id: str) -> CostRow | None:
|
||||
"""Poll /spend/logs for the call's row until it lands with a cost breakdown
|
||||
(rows flush ~60s behind the call via proxy_batch_write_at); None on timeout."""
|
||||
|
|
@ -156,6 +158,7 @@ def poll_cost_row(proxy: ProxyClient, request_id: str) -> CostRow | None:
|
|||
return None
|
||||
|
||||
|
||||
@step("Wait for a matching cost breakdown in the key's /spend/logs")
|
||||
def poll_cost_row_where(
|
||||
proxy: ProxyClient, api_key: str, predicate: Callable[[CostRow], bool]
|
||||
) -> CostRow | None:
|
||||
|
|
@ -182,6 +185,7 @@ def poll_cost_row_where(
|
|||
return None
|
||||
|
||||
|
||||
@step("Add a deployment with custom rates that calls {litellm_params.model}")
|
||||
def register_priced_model(
|
||||
proxy: ProxyClient,
|
||||
resources: ResourceManager,
|
||||
|
|
|
|||
|
|
@ -31,6 +31,7 @@ from e2e_http import (
|
|||
is_ok,
|
||||
unwrap,
|
||||
)
|
||||
from e2e_metadata import step
|
||||
from models import (
|
||||
AnthropicMessagesBody,
|
||||
ChatBody,
|
||||
|
|
@ -264,6 +265,7 @@ def _chat_body(
|
|||
class SpendClient:
|
||||
proxy: ProxyClient
|
||||
|
||||
@step('Send a /chat/completions request to {model} with the prompt "{content}"')
|
||||
def chat(
|
||||
self,
|
||||
key: str,
|
||||
|
|
@ -280,6 +282,7 @@ class SpendClient:
|
|||
_chat_body(model, content, max_tokens=max_tokens, tags=tags, user=user, cache=cache),
|
||||
)
|
||||
|
||||
@step('Send a streaming /chat/completions request to {model} with the prompt "{content}"')
|
||||
def chat_stream(
|
||||
self, key: str, model: str, content: str, *, max_tokens: int | None = None
|
||||
) -> StreamingResponse:
|
||||
|
|
@ -287,6 +290,7 @@ class SpendClient:
|
|||
key, _chat_body(model, content, max_tokens=max_tokens, stream=True)
|
||||
)
|
||||
|
||||
@step('Send a streaming /v1/messages request to {model} with the prompt "{content}"')
|
||||
def messages_stream(
|
||||
self, key: str, model: str, content: str, *, max_tokens: int
|
||||
) -> StreamingResponse:
|
||||
|
|
@ -300,9 +304,11 @@ class SpendClient:
|
|||
),
|
||||
)
|
||||
|
||||
@step('Send an /embeddings request to {model} for "{content}"')
|
||||
def embed(self, key: str, model: str, content: str) -> Result[EmbedResponse]:
|
||||
return self.proxy.embed(key, EmbedBody(model=model, input=content))
|
||||
|
||||
@step("Wait for at least {min_rows} of the key's spend logs in /spend/logs")
|
||||
def poll_logs_for_key(
|
||||
self,
|
||||
key: str,
|
||||
|
|
@ -314,6 +320,7 @@ class SpendClient:
|
|||
key, min_rows=min_rows, predicate=predicate
|
||||
)
|
||||
|
||||
@step('Estimate the cost of sending "{content}" to {model} with /spend/calculate')
|
||||
def calculate_spend(self, model: str, content: str) -> float:
|
||||
return unwrap(
|
||||
self.proxy.transport.post(
|
||||
|
|
@ -326,6 +333,7 @@ class SpendClient:
|
|||
)
|
||||
).cost
|
||||
|
||||
@step("Read the spend per tag from /spend/tags")
|
||||
def spend_by_tags(self) -> list[TagSpend]:
|
||||
result = self.proxy.transport.get(
|
||||
"/spend/tags",
|
||||
|
|
@ -339,6 +347,7 @@ class SpendClient:
|
|||
case _:
|
||||
return []
|
||||
|
||||
@step("Wait for the tag {tag} to reach the expected spend in /spend/tags")
|
||||
def poll_tag_spend(self, tag: str, *, minimum: float = 0.0) -> TagSpend | None:
|
||||
"""Poll /spend/tags until the tag's aggregate reaches `minimum`; last seen."""
|
||||
deadline = time.monotonic() + self.proxy.poll_timeout
|
||||
|
|
@ -354,6 +363,7 @@ class SpendClient:
|
|||
time.sleep(self.proxy.poll_interval)
|
||||
return entry
|
||||
|
||||
@step("Wait for the key's spend in /key/info to reach the expected minimum")
|
||||
def poll_key_spend(self, key: str, *, minimum: float = 0.0) -> float:
|
||||
deadline = time.monotonic() + self.proxy.poll_timeout
|
||||
spend = 0.0
|
||||
|
|
@ -364,6 +374,7 @@ class SpendClient:
|
|||
time.sleep(self.proxy.poll_interval)
|
||||
return spend
|
||||
|
||||
@step("Read the team's spend from /team/info")
|
||||
def team_spend(self, team_id: str) -> float:
|
||||
return (
|
||||
unwrap(
|
||||
|
|
@ -377,6 +388,7 @@ class SpendClient:
|
|||
or 0.0
|
||||
)
|
||||
|
||||
@step("Wait for the team's spend in /team/info to reach the expected minimum")
|
||||
def poll_team_spend(self, team_id: str, *, minimum: float = 0.0) -> float:
|
||||
outcome: Final = await_converged(
|
||||
lambda: self.team_spend(team_id),
|
||||
|
|
@ -388,6 +400,7 @@ class SpendClient:
|
|||
)
|
||||
return outcome.result if isinstance(outcome, Converged) else outcome.last_result
|
||||
|
||||
@step("Read the end user's spend from /customer/info")
|
||||
def customer_spend(self, customer_id: str) -> float:
|
||||
"""0.0 until the spend writer has upserted the end-user row, which /customer/info 404s before."""
|
||||
looked_up: Final = self.proxy.transport.get(
|
||||
|
|
@ -402,6 +415,7 @@ class SpendClient:
|
|||
case _:
|
||||
return 0.0
|
||||
|
||||
@step("Wait for the end user's spend in /customer/info to go above {minimum}")
|
||||
def poll_customer_spend(self, customer_id: str, *, minimum: float = 0.0) -> float:
|
||||
outcome: Final = await_converged(
|
||||
lambda: self.customer_spend(customer_id),
|
||||
|
|
@ -413,6 +427,7 @@ class SpendClient:
|
|||
)
|
||||
return outcome.result if isinstance(outcome, Converged) else outcome.last_result
|
||||
|
||||
@step("Scrape /metrics/ on every proxy replica")
|
||||
def scrape_metrics(self) -> Mapping[str, ProbeResult]:
|
||||
"""GET /metrics/ on every replica in PROXY_REPLICA_URLS, keyed by replica. The
|
||||
counter is per pod, so the union of the replicas is the fleet's exposition; the
|
||||
|
|
@ -425,6 +440,7 @@ class SpendClient:
|
|||
}
|
||||
)
|
||||
|
||||
@step("Read page {page} of /spend/logs/v2 at a page size of {page_size}")
|
||||
def spend_logs_page(
|
||||
self, *, api_key: str | None, page: int, page_size: int
|
||||
) -> SpendLogsPage:
|
||||
|
|
@ -447,9 +463,11 @@ class SpendClient:
|
|||
)
|
||||
)
|
||||
|
||||
@step("Call the management route {path}")
|
||||
def probe(self, path: str, *, params: DateRangeParams) -> ProbeResult:
|
||||
return self.proxy.transport.probe(path, params=params)
|
||||
|
||||
@step("Call the management route {path} until it answers successfully")
|
||||
def probe_until_healthy(self, path: str, *, params: DateRangeParams) -> ProbeResult:
|
||||
outcome: Final = await_converged(
|
||||
lambda: self.probe(path, params=params),
|
||||
|
|
@ -461,6 +479,7 @@ class SpendClient:
|
|||
)
|
||||
return outcome.result if isinstance(outcome, Converged) else outcome.last_result
|
||||
|
||||
@step("Create an internal user with the role {role}")
|
||||
def create_user(self, *, email: str, role: UserRole, user_id: str) -> str:
|
||||
return unwrap(
|
||||
self.proxy.transport.post(
|
||||
|
|
@ -471,6 +490,7 @@ class SpendClient:
|
|||
)
|
||||
).user_id
|
||||
|
||||
@step("Delete the internal user")
|
||||
def delete_user(self, user_id: str) -> None:
|
||||
_ = unwrap(
|
||||
self.proxy.transport.post(
|
||||
|
|
@ -481,6 +501,7 @@ class SpendClient:
|
|||
)
|
||||
)
|
||||
|
||||
@step("Generate a virtual key with {body}")
|
||||
def generate_key_record(self, body: KeyGenerateBody) -> KeyGenerateResponse:
|
||||
return unwrap(
|
||||
self.proxy.transport.post(
|
||||
|
|
@ -491,6 +512,7 @@ class SpendClient:
|
|||
)
|
||||
)
|
||||
|
||||
@step('Send a /chat/completions request to {model} with the prompt "{content}"')
|
||||
def send_chat(self, key: str, model: str, content: str, *, max_tokens: int) -> StreamingResponse:
|
||||
return self.proxy.transport.send(
|
||||
"/chat/completions",
|
||||
|
|
@ -498,6 +520,7 @@ class SpendClient:
|
|||
json=_chat_body(model, content, max_tokens=max_tokens),
|
||||
)
|
||||
|
||||
@step('Send a /queue/chat/completions request to {model} with the prompt "{content}"')
|
||||
def send_queued_chat(self, key: str, model: str, content: str, *, max_tokens: int) -> StreamingResponse:
|
||||
return self.proxy.transport.send(
|
||||
"/queue/chat/completions",
|
||||
|
|
@ -509,6 +532,7 @@ class SpendClient:
|
|||
),
|
||||
)
|
||||
|
||||
@step('Send a /v1/messages request to {model} with the prompt "{content}"')
|
||||
def send_messages(self, key: str, model: str, content: str, *, max_tokens: int) -> StreamingResponse:
|
||||
return self.proxy.transport.send(
|
||||
"/v1/messages",
|
||||
|
|
@ -520,9 +544,11 @@ class SpendClient:
|
|||
),
|
||||
)
|
||||
|
||||
@step('Send a /v1/responses request to {model} with the prompt "{content}"')
|
||||
def send_responses(self, key: str, model: str, content: str) -> StreamingResponse:
|
||||
return self.send_responses_with_headers(self.proxy.transport.bearer(key), model, content)
|
||||
|
||||
@step('Send a /v1/responses request to {model} with custom headers and the prompt "{content}"')
|
||||
def send_responses_with_headers(self, headers: AuthHeaders, model: str, content: str) -> StreamingResponse:
|
||||
return self.proxy.transport.send(
|
||||
"/v1/responses",
|
||||
|
|
@ -530,6 +556,7 @@ class SpendClient:
|
|||
json=ResponsesBody(model=model, input=content),
|
||||
)
|
||||
|
||||
@step('Send an /embeddings request to {model} for "{content}"')
|
||||
def send_embed(self, key: str, model: str, content: str) -> StreamingResponse:
|
||||
return self.proxy.transport.send(
|
||||
"/embeddings",
|
||||
|
|
@ -537,6 +564,7 @@ class SpendClient:
|
|||
json=EmbedBody(model=model, input=content),
|
||||
)
|
||||
|
||||
@step('Send a Gemini generateContent request to {model} through /gemini with the prompt "{content}"')
|
||||
def send_gemini_generate(self, key: str, model: str, content: str, *, max_tokens: int) -> StreamingResponse:
|
||||
return self.proxy.transport.send(
|
||||
f"/gemini/v1beta/models/{model}:generateContent",
|
||||
|
|
@ -547,6 +575,7 @@ class SpendClient:
|
|||
),
|
||||
)
|
||||
|
||||
@step("Upload a batch input file for {model} to /v1/files")
|
||||
def upload_batch_file(self, key: str, model: str, content: bytes) -> FileObject:
|
||||
return unwrap(
|
||||
self.proxy.transport.upload(
|
||||
|
|
@ -560,6 +589,7 @@ class SpendClient:
|
|||
)
|
||||
)
|
||||
|
||||
@step("Create a batch for {body.model} on /v1/batches")
|
||||
def create_batch(self, key: str, body: BatchCreateBody) -> BatchObject:
|
||||
return unwrap(
|
||||
self.proxy.transport.post(
|
||||
|
|
@ -570,6 +600,7 @@ class SpendClient:
|
|||
)
|
||||
)
|
||||
|
||||
@step("Retrieve the {provider} batch from /v1/batches")
|
||||
def retrieve_batch(self, key: str, batch_id: str, *, provider: str) -> BatchObject:
|
||||
return unwrap(
|
||||
self.proxy.transport.get(
|
||||
|
|
@ -580,6 +611,7 @@ class SpendClient:
|
|||
)
|
||||
)
|
||||
|
||||
@step("Post a callback log for {payload.model} to /v1/rust_control_plane/logs")
|
||||
def replay_callback_log(self, key: str, payload: CallbackLogPayload) -> CallbackLogsResponse:
|
||||
return unwrap(
|
||||
self.proxy.transport.post(
|
||||
|
|
@ -590,12 +622,15 @@ class SpendClient:
|
|||
)
|
||||
)
|
||||
|
||||
@step("Run a health check on {model} with /health")
|
||||
def health(self, model: str) -> ProbeResult:
|
||||
return self.proxy.transport.probe("/health", params=HealthParams(model=model))
|
||||
|
||||
@step("Read the key's daily activity from /user/daily/activity")
|
||||
def daily_activity_for_key(self, token: str, *, start: datetime, end: datetime) -> DailyActivityKeyBreakdown | None:
|
||||
return self._key_breakdown("/user/daily/activity", token, start=start, end=end)
|
||||
|
||||
@step("Read the key's usage export row from /user/daily/activity/aggregated")
|
||||
def usage_export_row_for_key(
|
||||
self, token: str, *, start: datetime, end: datetime
|
||||
) -> DailyActivityKeyBreakdown | None:
|
||||
|
|
@ -623,11 +658,13 @@ class SpendClient:
|
|||
None,
|
||||
)
|
||||
|
||||
@step("Wait for at least {min_requests} of the key's requests in /user/daily/activity")
|
||||
def poll_daily_activity_for_key(
|
||||
self, token: str, *, start: datetime, end: datetime, min_requests: int
|
||||
) -> DailyActivityKeyBreakdown | None:
|
||||
return self._poll_key_breakdown(lambda: self.daily_activity_for_key(token, start=start, end=end), min_requests)
|
||||
|
||||
@step("Wait for at least {min_requests} of the key's requests in /user/daily/activity/aggregated")
|
||||
def poll_usage_export_row_for_key(
|
||||
self, token: str, *, start: datetime, end: datetime, min_requests: int
|
||||
) -> DailyActivityKeyBreakdown | None:
|
||||
|
|
@ -648,6 +685,7 @@ class SpendClient:
|
|||
)
|
||||
return outcome.result if isinstance(outcome, Converged) else outcome.last_result
|
||||
|
||||
@step("Read the OpenAPI schema from /openapi.json")
|
||||
def openapi(self) -> OpenAPISchema:
|
||||
return unwrap(
|
||||
self.proxy.transport.get(
|
||||
|
|
|
|||
|
|
@ -5,6 +5,7 @@ from typing import Final
|
|||
|
||||
from e2e_config import provider_edge_base, unique_marker
|
||||
from e2e_http import unwrap
|
||||
from e2e_metadata import step
|
||||
from lifecycle import ResourceManager
|
||||
from models import ChatBody, ChatMessage, ChatResponse, KeyGenerateBody, LiteLLMParamsBody, TeamNewBody
|
||||
from spend_e2e_client import SpendClient
|
||||
|
|
@ -33,6 +34,10 @@ class TeamTraffic:
|
|||
return self.prompt_tokens * INPUT_RATE + self.completion_tokens * OUTPUT_RATE
|
||||
|
||||
|
||||
@step(
|
||||
"Add a priced deployment, then create two teams with one key each"
|
||||
" and send 7 /chat/completions requests per key, 6 of them at once"
|
||||
)
|
||||
def create_traffic(client: SpendClient, resources: ResourceManager) -> tuple[TeamTraffic, ...]:
|
||||
base: Final = provider_edge_base("openai")
|
||||
model: Final = f"e2e-reconciliation-{unique_marker()}"
|
||||
|
|
@ -85,6 +90,7 @@ def create_traffic(client: SpendClient, resources: ResourceManager) -> tuple[Tea
|
|||
return tuple(team_traffic() for _ in range(2))
|
||||
|
||||
|
||||
@step("Check that the /spend/logs rows of team {traffic.team_id} match each response's tokens and cost")
|
||||
def assert_logs_match(client: SpendClient, traffic: TeamTraffic) -> None:
|
||||
expected_ids: Final = frozenset(response.id for response in traffic.responses)
|
||||
assert len(expected_ids) == len(traffic.responses), "responses must have distinct IDs"
|
||||
|
|
|
|||
|
|
@ -36,7 +36,7 @@ from cost_rows import (
|
|||
)
|
||||
from e2e_config import CHEAP_OPENAI_MODEL, unique_marker
|
||||
from e2e_http import unwrap
|
||||
from e2e_metadata import Capability, Domain, Mode, Provider, Subject, meta
|
||||
from e2e_metadata import Capability, Domain, Mode, Provider, Route, Subject, meta
|
||||
from lifecycle import ResourceManager
|
||||
from models import (
|
||||
AnthropicMessagesBody,
|
||||
|
|
@ -184,6 +184,15 @@ class TestServiceTierPricing:
|
|||
assert_total_is_sum_of_components(row)
|
||||
|
||||
@pytest.mark.covers("quota_management.spend_tracking.service_tier_stream.records_served_tier")
|
||||
@meta(
|
||||
Subject(
|
||||
domain=Domain.SPEND_BUDGETS,
|
||||
route=Route.CHAT_COMPLETIONS,
|
||||
providers=(Provider.OPENAI,),
|
||||
models=(BACKEND,),
|
||||
mode=Mode.STREAM,
|
||||
)
|
||||
)
|
||||
def test_streamed_call_records_and_bills_the_served_tier(
|
||||
self, client: SpendClient, resources: ResourceManager, scoped_key: str
|
||||
) -> None:
|
||||
|
|
@ -232,6 +241,15 @@ class TestServiceTierPricing:
|
|||
assert_total_is_sum_of_components(row)
|
||||
|
||||
@pytest.mark.covers("llm.chat_completions.openai.service_tier.stream.echoes_served_tier")
|
||||
@meta(
|
||||
Subject(
|
||||
domain=Domain.LLM_TRANSLATION,
|
||||
route=Route.CHAT_COMPLETIONS,
|
||||
providers=(Provider.OPENAI,),
|
||||
models=(STREAM_BACKEND,),
|
||||
mode=Mode.STREAM,
|
||||
)
|
||||
)
|
||||
def test_every_streamed_chunk_carries_the_served_tier(
|
||||
self, client: SpendClient, resources: ResourceManager, scoped_key: str
|
||||
) -> None:
|
||||
|
|
@ -262,6 +280,15 @@ class TestServiceTierPricing:
|
|||
)
|
||||
|
||||
@pytest.mark.covers("quota_management.spend_tracking.service_tier_stream.responses_records_served_tier")
|
||||
@meta(
|
||||
Subject(
|
||||
domain=Domain.SPEND_BUDGETS,
|
||||
route=Route.RESPONSES,
|
||||
providers=(Provider.OPENAI,),
|
||||
models=(STREAM_BACKEND,),
|
||||
mode=Mode.STREAM,
|
||||
)
|
||||
)
|
||||
def test_responses_stream_records_the_served_tier(
|
||||
self, client: SpendClient, resources: ResourceManager, scoped_key: str
|
||||
) -> None:
|
||||
|
|
@ -308,6 +335,15 @@ class TestServiceTierPricing:
|
|||
assert_fresh_tokens_billed_at(row, INPUT_RATE_FOR_PRICING_BASIS[pricing_basis])
|
||||
|
||||
@pytest.mark.covers("quota_management.spend_tracking.service_tier_stream.messages_records_served_tier")
|
||||
@meta(
|
||||
Subject(
|
||||
domain=Domain.SPEND_BUDGETS,
|
||||
route=Route.MESSAGES,
|
||||
providers=(Provider.OPENAI,),
|
||||
models=(STREAM_BACKEND,),
|
||||
mode=Mode.STREAM,
|
||||
)
|
||||
)
|
||||
def test_messages_stream_records_the_served_tier(
|
||||
self, client: SpendClient, resources: ResourceManager, scoped_key: str
|
||||
) -> None:
|
||||
|
|
|
|||
|
|
@ -163,6 +163,12 @@ def test_schema_listed_spend_routes_are_responsive(client: SpendClient) -> None:
|
|||
assert not offenders, "non-responsive schema spend routes:\n" + "\n".join(offenders)
|
||||
|
||||
|
||||
@meta(
|
||||
Subject(
|
||||
domain=Domain.SPEND_BUDGETS,
|
||||
route=Route.SPEND_REPORTING,
|
||||
)
|
||||
)
|
||||
def test_capture_rate_reports_or_names_the_missing_billing_key(client: SpendClient) -> None:
|
||||
result: Final = client.probe(_CAPTURE_RATE_ROUTE, params=_date_range())
|
||||
print(f"{_CAPTURE_RATE_ROUTE} -> {result.status_code}\n{result.body[:600]}")
|
||||
|
|
|
|||
|
|
@ -32,6 +32,7 @@ from typing import Final
|
|||
import pytest
|
||||
from e2e_config import provider_edge_base, unique_marker
|
||||
from e2e_http import ProbeResult
|
||||
from e2e_metadata import Domain, Mode, Provider, Subject, meta
|
||||
from lifecycle import ResourceManager
|
||||
from models import ChatBody, ChatMessage, KeyGenerateBody, LiteLLMParamsBody, TeamNewBody
|
||||
from prometheus_client.parser import text_string_to_metric_families
|
||||
|
|
@ -41,6 +42,7 @@ from spend_reconciliation import INPUT_RATE, OUTPUT_RATE
|
|||
|
||||
pytestmark = pytest.mark.e2e
|
||||
|
||||
BACKEND: Final = "openai/gpt-5.6-luna"
|
||||
SPEND_METRIC: Final = "litellm_spend_metric_total"
|
||||
KEY_HASH_LABEL: Final = "hashed_api_key"
|
||||
TEAM_LABEL: Final = "team"
|
||||
|
|
@ -86,6 +88,14 @@ def _same_spend(actual: float | None, expected: float) -> bool:
|
|||
class TestSpendSurfaceConsistency:
|
||||
@pytest.mark.replayable
|
||||
@pytest.mark.covers("quota_management.spend_tracking.surface_consistency.matches_every_surface")
|
||||
@meta(
|
||||
Subject(
|
||||
domain=Domain.SPEND_BUDGETS,
|
||||
providers=(Provider.OPENAI,),
|
||||
models=(BACKEND,),
|
||||
mode=Mode.NONSTREAM,
|
||||
)
|
||||
)
|
||||
def test_one_request_lands_the_same_spend_on_every_surface(
|
||||
self, client: SpendClient, resources: ResourceManager
|
||||
) -> None:
|
||||
|
|
@ -96,7 +106,7 @@ class TestSpendSurfaceConsistency:
|
|||
model_id: Final = client.proxy.create_model(
|
||||
model,
|
||||
LiteLLMParamsBody(
|
||||
model="openai/gpt-5.6-luna",
|
||||
model=BACKEND,
|
||||
api_key="os.environ/OPENAI_API_KEY",
|
||||
api_base=None if base is None else f"{base}/v1",
|
||||
input_cost_per_token=INPUT_RATE,
|
||||
|
|
|
|||
|
|
@ -42,6 +42,7 @@ CLAUDE_MODEL = "claude-haiku-4-5"
|
|||
CODEX_MODEL = "openai-responses-codex"
|
||||
EMBEDDING_MODEL = "openai-text-embedding-3-small"
|
||||
OPENAI_BACKEND = "openai/gpt-5.5"
|
||||
ANTHROPIC_BACKEND: Final = "anthropic/claude-haiku-4-5"
|
||||
|
||||
|
||||
def _approx_equal(actual: float, expected: float) -> bool:
|
||||
|
|
@ -534,6 +535,15 @@ def test_end_user_spend_attributed_on_row(
|
|||
|
||||
|
||||
@pytest.mark.covers("quota_management.spend_tracking.end_user.attributes_responses_header")
|
||||
@meta(
|
||||
Subject(
|
||||
domain=Domain.SPEND_BUDGETS,
|
||||
route=Route.RESPONSES,
|
||||
providers=(Provider.OPENAI,),
|
||||
models=(CODEX_MODEL,),
|
||||
mode=Mode.NONSTREAM,
|
||||
)
|
||||
)
|
||||
@pytest.mark.parametrize("header", ["x-litellm-customer-id", "x-litellm-end-user-id"])
|
||||
def test_end_user_header_attributes_responses_row(
|
||||
client: SpendClient, scoped_key: str, resources: ResourceManager, header: str
|
||||
|
|
@ -659,6 +669,14 @@ def test_failure_call_writes_failure_status_row(
|
|||
|
||||
|
||||
@pytest.mark.covers("quota_management.spend_tracking.failure.writes_normalized_error")
|
||||
@meta(
|
||||
Subject(
|
||||
domain=Domain.SPEND_BUDGETS,
|
||||
providers=(Provider.ANTHROPIC, Provider.OPENAI),
|
||||
models=(OPENAI_BACKEND, ANTHROPIC_BACKEND),
|
||||
mode=Mode.NONSTREAM,
|
||||
)
|
||||
)
|
||||
def test_failure_rows_share_normalized_error_across_provider_wording(
|
||||
client: SpendClient, resources: ResourceManager, scoped_key: str
|
||||
) -> None:
|
||||
|
|
@ -668,7 +686,7 @@ def test_failure_rows_share_normalized_error_across_provider_wording(
|
|||
marker = unique_marker()
|
||||
deployments: Final = (
|
||||
(f"e2e-norm-openai-{marker}", OPENAI_BACKEND),
|
||||
(f"e2e-norm-anthropic-{marker}", "anthropic/claude-haiku-4-5"),
|
||||
(f"e2e-norm-anthropic-{marker}", ANTHROPIC_BACKEND),
|
||||
)
|
||||
for name, provider_model in deployments:
|
||||
model_id = client.proxy.create_model(
|
||||
|
|
@ -700,6 +718,14 @@ def test_failure_rows_share_normalized_error_across_provider_wording(
|
|||
|
||||
|
||||
@pytest.mark.covers("quota_management.spend_tracking.failure.attributes_provider")
|
||||
@meta(
|
||||
Subject(
|
||||
domain=Domain.SPEND_BUDGETS,
|
||||
providers=(Provider.OPENAI,),
|
||||
models=(OPENAI_BACKEND,),
|
||||
mode=Mode.NONSTREAM,
|
||||
)
|
||||
)
|
||||
def test_pre_call_rejection_row_attributes_provider_and_model_id(
|
||||
client: SpendClient, resources: ResourceManager
|
||||
) -> None:
|
||||
|
|
|
|||
|
|
@ -17,6 +17,7 @@ from typing import Final, Literal
|
|||
import pytest
|
||||
from e2e_config import unique_marker
|
||||
from e2e_http import unwrap
|
||||
from e2e_metadata import Capability, Domain, Mode, Provider, Route, Subject, meta
|
||||
from lifecycle import ResourceManager
|
||||
from models import (
|
||||
AnthropicContentBlock,
|
||||
|
|
@ -56,6 +57,16 @@ class TestWebSearchInterceptionSession:
|
|||
"quota_management.spend_tracking.websearch_interception.bills_under_request_session",
|
||||
exercised_on=("messages",),
|
||||
)
|
||||
@meta(
|
||||
Subject(
|
||||
domain=Domain.SPEND_BUDGETS,
|
||||
route=Route.MESSAGES,
|
||||
providers=(Provider.BEDROCK, Provider.PERPLEXITY),
|
||||
models=(BEDROCK_INVOKE_BACKEND,),
|
||||
capabilities=(Capability.WEB_SEARCH,),
|
||||
mode=Mode.NONSTREAM,
|
||||
)
|
||||
)
|
||||
def test_intercepted_search_is_billed_under_the_request_session(
|
||||
self, proxy: ProxyClient, resources: ResourceManager
|
||||
) -> None:
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue