diff --git a/tests/integration/README.md b/tests/integration/README.md index 1af3004b40e..c559e7545e0 100644 --- a/tests/integration/README.md +++ b/tests/integration/README.md @@ -30,8 +30,6 @@ Streaming checks send real HTTP transfer chunks, including one-byte partitions, The `messages_endpoint/` directory holds `/v1/messages` endpoint contracts: native-provider backends under `providers/` (`anthropic`, `bedrock`, `gemini`) and the translation bridges (`responses_bridge`, `chat_bridge`) at the top level. It runs in the providers shard; `run.py` selects test files recursively under each scheduled directory -A provider folder holds only what depends on that provider's wire format, and each subfolder is one feature that provider implements its own way: `headers/`, `streaming/`, `reasoning/`, `tools/`, `caching/`, `usage/` (reading the provider's token counts and pricing them), `multimodal/`, `context/` and `errors/`. A test goes in the folder of the feature it varies; one that fits no single folder tests two things and gets split. Behavior every provider shares, such as fallback or billing after a client disconnect, lives in the feature directory it exercises (`routing/`, `streaming/`, `spend/`). `_support/claude_code.py` holds a captured Claude Code request and stream builders that any directory can use as a realistic agent payload - The sdk shard exercises the SDK's own HTTP clients against local protocol peers with no gateway in the path, so a case here fails only when the client library or its wire behavior changes. The HTTP/2 case runs a hypercorn TLS peer offering h2 and http/1.1 over ALPN, drives the sync and async httpx handlers at it with `LITELLM_HTTP2` off and on, and asserts the version both the client and the peer observed on the wire. Put a test here only when it needs no proxy, database or Redis; a case that reaches the gateway belongs in one of the other shards The extensions shard uses the built-in generic callback and guardrail transports. It checks callback correlation and credential exclusion, guardrail rewriting and denial, retained OpenAI consumers and A2A wire versions. CircleCI runs it on parallel nodes, and each node starts its own database, Redis, upstream and proxy and runs its share of the group's files serially, split by recorded timings with `circleci tests split`. Tests keep the isolation of a serial run; they still must not assume a particular set of sibling files. `run.py --list` prints a group's files and `run.py ...` runs a subset of them diff --git a/tests/integration/_support/claude_code.py b/tests/integration/_support/claude_code.py index 909728ae5b0..044713e9656 100644 --- a/tests/integration/_support/claude_code.py +++ b/tests/integration/_support/claude_code.py @@ -9,7 +9,6 @@ from pydantic import JsonValue, TypeAdapter JSON_OBJECT: Final = TypeAdapter(dict[str, JsonValue]) ANTHROPIC_API_KEY: Final = "synthetic-anthropic-key" -SONNET: Final = "claude-sonnet-4-5" FABLE: Final = "claude-fable-5-1" OPUS: Final = "claude-opus-5-5" CLI_BETA: Final = ( diff --git a/tests/integration/messages_endpoint/providers/anthropic/caching/test_anthropic_prompt_cache_pricing_wire.py b/tests/integration/messages_endpoint/providers/anthropic/caching/test_anthropic_prompt_cache_pricing_wire.py deleted file mode 100644 index 06318b52aee..00000000000 --- a/tests/integration/messages_endpoint/providers/anthropic/caching/test_anthropic_prompt_cache_pricing_wire.py +++ /dev/null @@ -1,85 +0,0 @@ -import uuid -from typing import Final - -import pytest -from integration._support import claude_code as cc -from integration._support.client import Gateway, eventually -from integration._support.database import read_rows -from integration._support.wire import Reply, Request, wire_server - -_USAGE: Final = { - "input_tokens": 10, - "cache_read_input_tokens": 3000, - "cache_creation_input_tokens": 200, - "output_tokens": 5, -} - - -def test_cached_turn_charges_cache_read_and_creation_rates(gateway: Gateway) -> None: - identity: Final = f"msg_pc_{uuid.uuid4().hex}" - turn1: Final = cc.frontier_request( - f"cache-bust-{uuid.uuid4().hex}", - "high", - 64000, - prompt_text="Read /tmp/cc_probe/hello.txt and reply with its single word", - ) - turn2: Final = cc.tool_loop_turn2( - turn1, - ( - {"type": "thinking", "thinking": "need to read the file", "signature": "sig_anthropic_1"}, - { - "type": "tool_use", - "id": "toolu_read_1", - "name": "Read", - "input": {"file_path": "/tmp/cc_probe/hello.txt"}, - }, - ), - (("toolu_read_1", "1\tPROBE\n2\t"),), - ) - - def respond(request: Request) -> Reply: - assert request.method == "POST" - assert request.target == "/v1/messages", request.target - body: Final = cc.JSON_OBJECT.validate_json(request.body) - expected: Final = {**turn2, "model": cc.FABLE} - assert body == expected, { - key: (expected.get(key), body.get(key)) - for key in expected.keys() | body.keys() - if expected.get(key) != body.get(key) - } - return Reply(content_type="text/event-stream", chunks=cc.text_stream(identity, cc.FABLE, "PROBE", _USAGE)) - - with wire_server(respond) as wire, gateway.scenario() as scenario: - model: Final = scenario.model( - model=f"anthropic/{cc.FABLE}", - api_base=wire.url, - api_key=cc.ANTHROPIC_API_KEY, - input_cost_per_token=1e-6, - output_cost_per_token=5e-6, - cache_read_input_token_cost=1e-7, - cache_creation_input_token_cost=1.25e-6, - ) - response: Final = gateway.request( - "POST", - "/v1/messages", - {**turn2, "model": model}, - params={"beta": "true"}, - headers=cc.cli_headers(gateway.key, cc.FRONTIER_CLI_BETA), - ) - assert response.status_code == 200, response.text - events: Final = cc.sse_events(response.text) - usage: Final = events[0][1]["message"]["usage"] - assert usage["cache_read_input_tokens"] == 3000, usage - assert usage["cache_creation_input_tokens"] == 200, usage - assert len(wire.drain()) == 1 - rows: Final = eventually( - lambda: read_rows( - 'SELECT spend, prompt_tokens, completion_tokens FROM "LiteLLM_SpendLogs" WHERE request_id=%s', - (identity,), - ), - lambda values: len(values) == 1, - seconds=70, - ) - assert float(rows[0]["spend"]) == pytest.approx(10 * 1e-6 + 3000 * 1e-7 + 200 * 1.25e-6 + 5 * 5e-6), dict( - rows[0] - ) diff --git a/tests/integration/messages_endpoint/providers/anthropic/context/test_anthropic_compaction_wire.py b/tests/integration/messages_endpoint/providers/anthropic/context/test_anthropic_compaction_wire.py deleted file mode 100644 index c6ac00c20e8..00000000000 --- a/tests/integration/messages_endpoint/providers/anthropic/context/test_anthropic_compaction_wire.py +++ /dev/null @@ -1,108 +0,0 @@ -import uuid -from typing import Final - -from integration._support import claude_code as cc -from integration._support.client import Gateway -from integration._support.wire import Reply, Request, wire_server - -_CONTEXT_MANAGEMENT: Final = { - "edits": [ - {"type": "clear_thinking_20251015", "keep": "all"}, - {"type": "compact_20260112", "trigger": {"type": "input_tokens", "value": 150000}}, - ] -} -_COMPACTION_BLOCK: Final = {"type": "compaction", "content": ""} - - -def _compaction_stream(identity: str) -> tuple[bytes, ...]: - return ( - cc.sse_frame( - "message_start", - { - "type": "message_start", - "message": { - "id": identity, - "type": "message", - "role": "assistant", - "model": cc.FABLE, - "content": [], - "stop_reason": None, - "stop_sequence": None, - "usage": {"input_tokens": 20, "output_tokens": 1}, - }, - }, - ), - cc.sse_frame( - "content_block_start", - {"type": "content_block_start", "index": 0, "content_block": dict(_COMPACTION_BLOCK)}, - ), - cc.sse_frame("content_block_stop", {"type": "content_block_stop", "index": 0}), - cc.sse_frame( - "message_delta", - { - "type": "message_delta", - "delta": {"stop_reason": "end_turn", "stop_sequence": None}, - "usage": {"output_tokens": 5}, - "context_management": { - "applied_edits": [{"type": "compact_20260112", "compacted_at": "2026-09-26T00:00:00Z"}] - }, - }, - ), - cc.sse_frame("message_stop", {"type": "message_stop"}), - ) - - -def test_compaction_edit_and_applied_edit_block_round_trip_through_anthropic(gateway: Gateway) -> None: - identity: Final = f"msg_cm_{uuid.uuid4().hex}" - request_body: Final = { - **cc.frontier_request(f"cache-bust-{uuid.uuid4().hex}", "high", 64000), - "context_management": _CONTEXT_MANAGEMENT, - } - turn3: Final = cc.tool_loop_turn2( - request_body, - ( - dict(_COMPACTION_BLOCK), - {"type": "text", "text": "continuing after compaction"}, - ), - (), - ) - first_expected: Final = {**request_body, "model": cc.FABLE} - second_expected: Final = {**turn3, "model": cc.FABLE} - - def respond(request: Request) -> Reply: - assert request.method == "POST" - assert request.target == "/v1/messages", request.target - body: Final = cc.JSON_OBJECT.validate_json(request.body) - if body == first_expected: - assert body["context_management"] == _CONTEXT_MANAGEMENT - return Reply(content_type="text/event-stream", chunks=_compaction_stream(identity)) - assert body == second_expected, { - key: (second_expected.get(key), body.get(key)) - for key in second_expected.keys() | body.keys() - if second_expected.get(key) != body.get(key) - } - assert body["context_management"] == _CONTEXT_MANAGEMENT - return Reply( - content_type="text/event-stream", - chunks=cc.text_stream("msg_cm_next", cc.FABLE, "OK", {"input_tokens": 20, "output_tokens": 2}), - ) - - with wire_server(respond) as wire, gateway.scenario() as scenario: - model: Final = scenario.model(model=f"anthropic/{cc.FABLE}", api_base=wire.url, api_key=cc.ANTHROPIC_API_KEY) - headers: Final = cc.cli_headers(gateway.key, cc.FRONTIER_CLI_BETA) - response1: Final = gateway.request( - "POST", "/v1/messages", {**request_body, "model": model}, params={"beta": "true"}, headers=headers - ) - assert response1.status_code == 200, response1.text - events: Final = cc.sse_events(response1.text) - assert events[1][1]["content_block"] == _COMPACTION_BLOCK, events[1] - deltas: Final = [data for event, data in events if event == "message_delta"] - assert len(deltas) == 1 and deltas[0].get("context_management") == { - "applied_edits": [{"type": "compact_20260112", "compacted_at": "2026-09-26T00:00:00Z"}] - }, events - response2: Final = gateway.request( - "POST", "/v1/messages", {**turn3, "model": model}, params={"beta": "true"}, headers=headers - ) - assert response2.status_code == 200, response2.text - bodies: Final = tuple(cc.JSON_OBJECT.validate_json(request.body) for request in wire.drain()) - assert bodies == (first_expected, second_expected), bodies diff --git a/tests/integration/messages_endpoint/providers/anthropic/errors/test_anthropic_bare_string_content_rejected_wire.py b/tests/integration/messages_endpoint/providers/anthropic/errors/test_anthropic_bare_string_content_rejected_wire.py deleted file mode 100644 index ffea1830e54..00000000000 --- a/tests/integration/messages_endpoint/providers/anthropic/errors/test_anthropic_bare_string_content_rejected_wire.py +++ /dev/null @@ -1,32 +0,0 @@ -import json -import time -import uuid -from typing import Final - -import pytest -from integration._support.client import Gateway, eventually, object_value -from integration._support.database import read_rows -from integration._support.wire import Reply, Request, wire_server - - -@pytest.mark.covers("other.provider_wire.anthropic.bare_string_content_item_is_client_error") -@pytest.mark.parametrize( - "text", [pytest.param("what type of file is this?", id="type_word"), pytest.param("hello", id="plain")] -) -def test_anthropic_bare_string_content_item_is_rejected_as_client_error_before_the_wire( - gateway: Gateway, text: str -) -> None: - def respond(request: Request) -> Reply: - raise AssertionError(f"upstream must not be reached: {request.target}") - - with wire_server(respond) as wire, gateway.scenario() as scenario: - model: Final = scenario.model( - model="anthropic/claude-sonnet-4-5-20250929", api_base=wire.url, api_key="synthetic-anthropic-key" - ) - response: Final = gateway.request( - "POST", - "/v1/chat/completions", - {"model": model, "max_tokens": 16, "timeout": 5, "messages": [{"role": "system", "content": [text]}]}, - ) - assert response.status_code == 400, response.text - assert wire.drain() == () diff --git a/tests/integration/messages_endpoint/providers/anthropic/errors/test_anthropic_slow_upstream_cutoff_wire.py b/tests/integration/messages_endpoint/providers/anthropic/errors/test_anthropic_slow_upstream_cutoff_wire.py deleted file mode 100644 index 7532afe718d..00000000000 --- a/tests/integration/messages_endpoint/providers/anthropic/errors/test_anthropic_slow_upstream_cutoff_wire.py +++ /dev/null @@ -1,64 +0,0 @@ -import json -import time -import uuid -from typing import Final - -import pytest -from integration._support.client import Gateway, eventually, object_value -from integration._support.database import read_rows -from integration._support.wire import Reply, Request, wire_server - - -@pytest.mark.covers("other.provider_wire.anthropic.messages_request_timeout_reaches_transport") -def test_anthropic_messages_slow_upstream_is_cut_off_at_the_deployment_request_timeout(gateway: Gateway) -> None: - identity: Final = "anthropic-timeout-" + uuid.uuid4().hex - prompt: Final = f"slow answer {identity}" - - def respond(request: Request) -> Reply: - assert request.method == "POST" and request.target == "/v1/messages" - assert request.headers["x-api-key"] == "synthetic-anthropic-key" - body: Final = json.loads(request.body) - assert body["model"] == "claude-sonnet-4-5-20250929" - assert body["max_tokens"] == 16 - assert body["messages"] == [{"role": "user", "content": prompt}] - assert not { - "timeout", - "request_timeout", - "stream_chunk_size", - "litellm_params", - "litellm_metadata", - "rpm", - "tpm", - }.intersection(body) - time.sleep(1.5) - return Reply( - body=json.dumps( - { - "id": identity, - "type": "message", - "role": "assistant", - "model": "claude-sonnet-4-5-20250929", - "content": [{"type": "text", "text": "late"}], - "stop_reason": "end_turn", - "stop_sequence": None, - "usage": {"input_tokens": 3, "output_tokens": 1}, - } - ).encode() - ) - - with wire_server(respond) as wire, gateway.scenario() as scenario: - model: Final = scenario.model( - model="anthropic/claude-sonnet-4-5-20250929", - api_base=wire.url, - api_key="synthetic-anthropic-key", - request_timeout=0.3, - ) - response: Final = gateway.request( - "POST", - "/v1/messages", - {"model": model, "max_tokens": 16, "messages": [{"role": "user", "content": prompt}]}, - headers={"anthropic-version": "2023-06-01"}, - ) - assert response.status_code == 408, response.text - assert "Timeout" in response.json()["error"]["message"], response.text - assert eventually(wire.drain, lambda requests: len(requests) == 1, seconds=5, return_last_on_timeout=True) diff --git a/tests/integration/messages_endpoint/providers/anthropic/headers/test_anthropic_request_fidelity_wire.py b/tests/integration/messages_endpoint/providers/anthropic/headers/test_anthropic_request_fidelity_wire.py deleted file mode 100644 index a5a57737fc1..00000000000 --- a/tests/integration/messages_endpoint/providers/anthropic/headers/test_anthropic_request_fidelity_wire.py +++ /dev/null @@ -1,71 +0,0 @@ -import uuid -from typing import Final - -from integration._support import claude_code as cc -from integration._support.client import Gateway, eventually -from integration._support.database import read_rows -from integration._support.wire import Reply, Request, wire_server - -_MODEL: Final = cc.SONNET - - -def test_streaming_request_reaches_anthropic_intact_and_streams_back(gateway: Gateway) -> None: - identity: Final = f"msg_cc_{uuid.uuid4().hex}" - request_body: Final = cc.claude_code_request(f"cache-bust-{uuid.uuid4().hex}") - cli_beta: Final = frozenset(cc.CLI_BETA.split(",")) - - def respond(request: Request) -> Reply: - assert request.method == "POST" - assert request.target == "/v1/messages", request.target - assert request.headers["x-api-key"] == cc.ANTHROPIC_API_KEY - assert request.headers["anthropic-version"] == "2023-06-01" - assert frozenset(request.headers.get("anthropic-beta", "").split(",")) == cli_beta, request.headers.get( - "anthropic-beta" - ) - assert "authorization" not in request.headers, dict(request.headers) - assert all(gateway.key not in value for value in request.headers.values()), dict(request.headers) - body: Final = cc.JSON_OBJECT.validate_json(request.body) - expected: Final = {**request_body, "model": _MODEL} - assert body == expected, { - key: (expected.get(key), body.get(key)) - for key in expected.keys() | body.keys() - if expected.get(key) != body.get(key) - } - return Reply( - content_type="text/event-stream", - chunks=cc.text_stream(identity, _MODEL, "PONG", {"input_tokens": 12, "output_tokens": 4}), - ) - - with wire_server(respond) as wire, gateway.scenario() as scenario: - model: Final = scenario.model(model=f"anthropic/{_MODEL}", api_base=wire.url, api_key=cc.ANTHROPIC_API_KEY) - response: Final = gateway.request( - "POST", - "/v1/messages", - {**request_body, "model": model}, - params={"beta": "true"}, - headers=cc.cli_headers(gateway.key), - ) - assert response.status_code == 200, response.text - assert response.headers["content-type"].startswith("text/event-stream"), dict(response.headers) - events: Final = cc.sse_events(response.text) - assert [event for event, _ in events] == [ - "message_start", - "content_block_start", - "content_block_delta", - "content_block_stop", - "message_delta", - "message_stop", - ] - assert events[2][1]["delta"] == {"type": "text_delta", "text": "PONG"} - assert events[4][1]["delta"]["stop_reason"] == "end_turn" - assert events[4][1]["usage"]["output_tokens"] == 4 - assert len(wire.drain()) == 1 - rows: Final = eventually( - lambda: read_rows( - 'SELECT prompt_tokens, completion_tokens FROM "LiteLLM_SpendLogs" WHERE request_id=%s', - (identity,), - ), - lambda values: len(values) == 1, - seconds=70, - ) - assert rows[0]["prompt_tokens"] == 12 and rows[0]["completion_tokens"] == 4 diff --git a/tests/integration/messages_endpoint/providers/anthropic/multimodal/test_anthropic_document_input_wire.py b/tests/integration/messages_endpoint/providers/anthropic/multimodal/test_anthropic_document_input_wire.py deleted file mode 100644 index 48b914a7374..00000000000 --- a/tests/integration/messages_endpoint/providers/anthropic/multimodal/test_anthropic_document_input_wire.py +++ /dev/null @@ -1,133 +0,0 @@ -import base64 -import uuid -from typing import Final - -from integration._support import claude_code as cc -from integration._support.client import Gateway -from integration._support.wire import Reply, Request, wire_server - -_PDF_BYTES: Final = ( - b"%PDF-1.1\n" - b"1 0 obj<>endobj\n" - b"2 0 obj<>endobj\n" - b"3 0 obj<>endobj\n" - b"trailer<>\n%%EOF" -) -_DOC_BLOCK: Final = { - "type": "document", - "source": {"type": "base64", "data": base64.b64encode(_PDF_BYTES).decode(), "media_type": "application/pdf"}, - "citations": {"enabled": True}, -} - - -def _cited_stream(identity: str) -> tuple[bytes, ...]: - return ( - cc.sse_frame( - "message_start", - { - "type": "message_start", - "message": { - "id": identity, - "type": "message", - "role": "assistant", - "model": cc.FABLE, - "content": [], - "stop_reason": None, - "stop_sequence": None, - "usage": {"input_tokens": 20, "output_tokens": 1}, - }, - }, - ), - cc.sse_frame( - "content_block_start", - {"type": "content_block_start", "index": 0, "content_block": {"type": "text", "text": ""}}, - ), - cc.sse_frame( - "content_block_delta", - {"type": "content_block_delta", "index": 0, "delta": {"type": "text_delta", "text": "A page."}}, - ), - cc.sse_frame( - "content_block_delta", - { - "type": "content_block_delta", - "index": 0, - "delta": { - "type": "citations_delta", - "citation": { - "type": "page_location", - "document_index": 0, - "document_title": "dot.pdf", - "start_page_number": 1, - "end_page_number": 1, - "cited_text": "Page", - }, - }, - }, - ), - cc.sse_frame("content_block_stop", {"type": "content_block_stop", "index": 0}), - cc.sse_frame( - "message_delta", - { - "type": "message_delta", - "delta": {"stop_reason": "end_turn", "stop_sequence": None}, - "usage": {"output_tokens": 6}, - }, - ), - cc.sse_frame("message_stop", {"type": "message_stop"}), - ) - - -def test_base64_pdf_document_with_citations_reaches_anthropic_identical(gateway: Gateway) -> None: - request_body: Final = { - **cc.frontier_request(f"cache-bust-{uuid.uuid4().hex}", "high", 64000), - "messages": [ - { - "role": "user", - "content": [ - dict(_DOC_BLOCK), - {"type": "text", "text": f"What is on page one? {uuid.uuid4().hex}"}, - ], - } - ], - } - - def respond(request: Request) -> Reply: - assert request.method == "POST" - assert request.target == "/v1/messages", request.target - body: Final = cc.JSON_OBJECT.validate_json(request.body) - expected: Final = {**request_body, "model": cc.FABLE} - assert body == expected, { - key: (expected.get(key), body.get(key)) - for key in expected.keys() | body.keys() - if expected.get(key) != body.get(key) - } - return Reply(content_type="text/event-stream", chunks=_cited_stream(f"msg_doc_{uuid.uuid4().hex}")) - - with wire_server(respond) as wire, gateway.scenario() as scenario: - model: Final = scenario.model(model=f"anthropic/{cc.FABLE}", api_base=wire.url, api_key=cc.ANTHROPIC_API_KEY) - response: Final = gateway.request( - "POST", - "/v1/messages", - {**request_body, "model": model}, - params={"beta": "true"}, - headers=cc.cli_headers(gateway.key, cc.FRONTIER_CLI_BETA), - ) - assert response.status_code == 200, response.text - events: Final = cc.sse_events(response.text) - citations: Final = [ - data["delta"] for event, data in events if data.get("delta", {}).get("type") == "citations_delta" - ] - assert citations == [ - { - "type": "citations_delta", - "citation": { - "type": "page_location", - "document_index": 0, - "document_title": "dot.pdf", - "start_page_number": 1, - "end_page_number": 1, - "cited_text": "Page", - }, - } - ], citations - assert len(wire.drain()) == 1 diff --git a/tests/integration/messages_endpoint/providers/anthropic/multimodal/test_anthropic_image_input_wire.py b/tests/integration/messages_endpoint/providers/anthropic/multimodal/test_anthropic_image_input_wire.py deleted file mode 100644 index 8a929ee2a06..00000000000 --- a/tests/integration/messages_endpoint/providers/anthropic/multimodal/test_anthropic_image_input_wire.py +++ /dev/null @@ -1,67 +0,0 @@ -import uuid -from typing import Final - -from integration._support import claude_code as cc -from integration._support.client import Gateway -from integration._support.wire import Reply, Request, wire_server - -_PNG_B64: Final = "iVBORw0KGgoAAAANSUhEUgAAAAQAAAAECAIAAAAmkwkpAAAAEElEQVR4nGP4z8AARwzEcQCukw/x0F8jngAAAABJRU5ErkJggg==" -_IMAGE_BLOCK: Final = { - "type": "image", - "source": {"type": "base64", "data": _PNG_B64, "media_type": "image/png"}, -} - - -def test_tool_result_image_block_and_pasted_image_reach_anthropic_identical(gateway: Gateway) -> None: - turn1: Final = cc.frontier_request( - f"cache-bust-{uuid.uuid4().hex}", - "high", - 64000, - prompt_text="Read /tmp/cc_probe/dot.png and say what colour it is", - ) - with_image_result: Final = cc.tool_loop_turn2( - turn1, - ({"type": "tool_use", "id": "toolu_img", "name": "Read", "input": {"file_path": "/tmp/cc_probe/dot.png"}},), - (("toolu_img", [dict(_IMAGE_BLOCK)]),), - ) - pasted: Final = { - **cc.frontier_request(f"cache-bust-{uuid.uuid4().hex}", "high", 64000), - "messages": [ - { - "role": "user", - "content": [ - dict(_IMAGE_BLOCK), - {"type": "text", "text": f"What colour is this? {uuid.uuid4().hex}"}, - ], - } - ], - } - first_expected: Final = {**with_image_result, "model": cc.FABLE} - second_expected: Final = {**pasted, "model": cc.FABLE} - - def respond(request: Request) -> Reply: - assert request.method == "POST" - assert request.target == "/v1/messages", request.target - body: Final = cc.JSON_OBJECT.validate_json(request.body) - if body != first_expected: - assert body == second_expected, body - return Reply( - content_type="text/event-stream", - chunks=cc.text_stream( - f"msg_img_{uuid.uuid4().hex}", cc.FABLE, "RED", {"input_tokens": 20, "output_tokens": 2} - ), - ) - - with wire_server(respond) as wire, gateway.scenario() as scenario: - model: Final = scenario.model(model=f"anthropic/{cc.FABLE}", api_base=wire.url, api_key=cc.ANTHROPIC_API_KEY) - headers: Final = cc.cli_headers(gateway.key, cc.FRONTIER_CLI_BETA) - response1: Final = gateway.request( - "POST", "/v1/messages", {**with_image_result, "model": model}, params={"beta": "true"}, headers=headers - ) - assert response1.status_code == 200, response1.text - response2: Final = gateway.request( - "POST", "/v1/messages", {**pasted, "model": model}, params={"beta": "true"}, headers=headers - ) - assert response2.status_code == 200, response2.text - bodies: Final = tuple(cc.JSON_OBJECT.validate_json(request.body) for request in wire.drain()) - assert bodies == (first_expected, second_expected), bodies diff --git a/tests/integration/messages_endpoint/providers/anthropic/tools/test_anthropic_advisor_wire.py b/tests/integration/messages_endpoint/providers/anthropic/test_anthropic_advisor_wire.py similarity index 100% rename from tests/integration/messages_endpoint/providers/anthropic/tools/test_anthropic_advisor_wire.py rename to tests/integration/messages_endpoint/providers/anthropic/test_anthropic_advisor_wire.py diff --git a/tests/integration/messages_endpoint/providers/anthropic/reasoning/test_anthropic_legacy_thinking_budget_wire.py b/tests/integration/messages_endpoint/providers/anthropic/test_anthropic_legacy_thinking_budget_wire.py similarity index 100% rename from tests/integration/messages_endpoint/providers/anthropic/reasoning/test_anthropic_legacy_thinking_budget_wire.py rename to tests/integration/messages_endpoint/providers/anthropic/test_anthropic_legacy_thinking_budget_wire.py diff --git a/tests/integration/messages_endpoint/providers/anthropic/streaming/test_anthropic_messages_live_lifecycle_wire.py b/tests/integration/messages_endpoint/providers/anthropic/test_anthropic_messages_live_lifecycle_wire.py similarity index 100% rename from tests/integration/messages_endpoint/providers/anthropic/streaming/test_anthropic_messages_live_lifecycle_wire.py rename to tests/integration/messages_endpoint/providers/anthropic/test_anthropic_messages_live_lifecycle_wire.py diff --git a/tests/integration/messages_endpoint/providers/anthropic/errors/test_anthropic_messages_timeout_wire.py b/tests/integration/messages_endpoint/providers/anthropic/test_anthropic_messages_timeout_wire.py similarity index 100% rename from tests/integration/messages_endpoint/providers/anthropic/errors/test_anthropic_messages_timeout_wire.py rename to tests/integration/messages_endpoint/providers/anthropic/test_anthropic_messages_timeout_wire.py diff --git a/tests/integration/messages_endpoint/providers/anthropic/reasoning/test_anthropic_thinking_signature_retry_wire.py b/tests/integration/messages_endpoint/providers/anthropic/test_anthropic_thinking_signature_retry_wire.py similarity index 100% rename from tests/integration/messages_endpoint/providers/anthropic/reasoning/test_anthropic_thinking_signature_retry_wire.py rename to tests/integration/messages_endpoint/providers/anthropic/test_anthropic_thinking_signature_retry_wire.py diff --git a/tests/integration/messages_endpoint/providers/anthropic/caching/test_anthropic_tool_history_cache_tokens_wire.py b/tests/integration/messages_endpoint/providers/anthropic/test_anthropic_wire.py similarity index 63% rename from tests/integration/messages_endpoint/providers/anthropic/caching/test_anthropic_tool_history_cache_tokens_wire.py rename to tests/integration/messages_endpoint/providers/anthropic/test_anthropic_wire.py index eff5fd565b5..7e9c5be227a 100644 --- a/tests/integration/messages_endpoint/providers/anthropic/caching/test_anthropic_tool_history_cache_tokens_wire.py +++ b/tests/integration/messages_endpoint/providers/anthropic/test_anthropic_wire.py @@ -126,3 +126,81 @@ def test_anthropic_tool_history_and_cache_tokens_keep_wire_and_accounting_contra parsed: Final = json.loads(metadata) if isinstance(metadata, str) else object_value(metadata) assert parsed["cost_breakdown"]["input_cost"] == pytest.approx(0.0245) assert parsed["cost_breakdown"]["output_cost"] == pytest.approx(0.008) + + +@pytest.mark.covers("other.provider_wire.anthropic.bare_string_content_item_is_client_error") +@pytest.mark.parametrize( + "text", [pytest.param("what type of file is this?", id="type_word"), pytest.param("hello", id="plain")] +) +def test_anthropic_bare_string_content_item_is_rejected_as_client_error_before_the_wire( + gateway: Gateway, text: str +) -> None: + def respond(request: Request) -> Reply: + raise AssertionError(f"upstream must not be reached: {request.target}") + + with wire_server(respond) as wire, gateway.scenario() as scenario: + model: Final = scenario.model( + model="anthropic/claude-sonnet-4-5-20250929", api_base=wire.url, api_key="synthetic-anthropic-key" + ) + response: Final = gateway.request( + "POST", + "/v1/chat/completions", + {"model": model, "max_tokens": 16, "timeout": 5, "messages": [{"role": "system", "content": [text]}]}, + ) + assert response.status_code == 400, response.text + assert wire.drain() == () + + +@pytest.mark.covers("other.provider_wire.anthropic.messages_request_timeout_reaches_transport") +def test_anthropic_messages_slow_upstream_is_cut_off_at_the_deployment_request_timeout(gateway: Gateway) -> None: + identity: Final = "anthropic-timeout-" + uuid.uuid4().hex + prompt: Final = f"slow answer {identity}" + + def respond(request: Request) -> Reply: + assert request.method == "POST" and request.target == "/v1/messages" + assert request.headers["x-api-key"] == "synthetic-anthropic-key" + body: Final = json.loads(request.body) + assert body["model"] == "claude-sonnet-4-5-20250929" + assert body["max_tokens"] == 16 + assert body["messages"] == [{"role": "user", "content": prompt}] + assert not { + "timeout", + "request_timeout", + "stream_chunk_size", + "litellm_params", + "litellm_metadata", + "rpm", + "tpm", + }.intersection(body) + time.sleep(1.5) + return Reply( + body=json.dumps( + { + "id": identity, + "type": "message", + "role": "assistant", + "model": "claude-sonnet-4-5-20250929", + "content": [{"type": "text", "text": "late"}], + "stop_reason": "end_turn", + "stop_sequence": None, + "usage": {"input_tokens": 3, "output_tokens": 1}, + } + ).encode() + ) + + with wire_server(respond) as wire, gateway.scenario() as scenario: + model: Final = scenario.model( + model="anthropic/claude-sonnet-4-5-20250929", + api_base=wire.url, + api_key="synthetic-anthropic-key", + request_timeout=0.3, + ) + response: Final = gateway.request( + "POST", + "/v1/messages", + {"model": model, "max_tokens": 16, "messages": [{"role": "user", "content": prompt}]}, + headers={"anthropic-version": "2023-06-01"}, + ) + assert response.status_code == 408, response.text + assert "Timeout" in response.json()["error"]["message"], response.text + assert eventually(wire.drain, lambda requests: len(requests) == 1, seconds=5, return_last_on_timeout=True) diff --git a/tests/integration/messages_endpoint/providers/anthropic/tools/test_websearch_interception_wire.py b/tests/integration/messages_endpoint/providers/anthropic/test_websearch_interception_wire.py similarity index 100% rename from tests/integration/messages_endpoint/providers/anthropic/tools/test_websearch_interception_wire.py rename to tests/integration/messages_endpoint/providers/anthropic/test_websearch_interception_wire.py diff --git a/tests/integration/messages_endpoint/providers/anthropic/tools/test_anthropic_tool_loop_wire.py b/tests/integration/messages_endpoint/providers/anthropic/tools/test_anthropic_tool_loop_wire.py deleted file mode 100644 index f3fb9786e02..00000000000 --- a/tests/integration/messages_endpoint/providers/anthropic/tools/test_anthropic_tool_loop_wire.py +++ /dev/null @@ -1,165 +0,0 @@ -import uuid -from typing import Final - -from integration._support import claude_code as cc -from integration._support.client import Gateway, eventually -from integration._support.database import read_rows -from integration._support.wire import Reply, Request, wire_server -from pydantic import JsonValue - -_THINKING: Final = "need to read the file" -_SIGNATURE: Final = "sig_probe_1" -_USAGE: Final = {"input_tokens": 20, "output_tokens": 10} - - -def _diff(expected: dict[str, JsonValue], body: dict[str, JsonValue]) -> dict[str, JsonValue]: - return { - key: {"expected": expected.get(key), "upstream": body.get(key)} - for key in expected.keys() | body.keys() - if expected.get(key) != body.get(key) - } - - -def test_tool_loop_round_trips_thinking_tool_use_and_tool_result(gateway: Gateway) -> None: - identity1: Final = f"msg_tl1_{uuid.uuid4().hex}" - identity2: Final = f"msg_tl2_{uuid.uuid4().hex}" - turn1: Final = cc.frontier_request( - f"cache-bust-{uuid.uuid4().hex}", - "high", - 64000, - prompt_text="Read /tmp/cc_probe/hello.txt and reply with its single word", - ) - calls: Final = (("toolu_read_1", "Read", {"file_path": "/tmp/cc_probe/hello.txt"}),) - turn2: Final = cc.tool_loop_turn2( - turn1, - ( - {"type": "thinking", "thinking": _THINKING, "signature": _SIGNATURE}, - { - "type": "tool_use", - "id": "toolu_read_1", - "name": "Read", - "input": {"file_path": "/tmp/cc_probe/hello.txt"}, - }, - ), - (("toolu_read_1", "1\tPROBE\n2\t"),), - ) - first_expected: Final = {**turn1, "model": cc.FABLE} - second_expected: Final = {**turn2, "model": cc.FABLE} - - def respond(request: Request) -> Reply: - assert request.method == "POST" - assert request.target == "/v1/messages", request.target - body: Final = cc.JSON_OBJECT.validate_json(request.body) - if body == first_expected: - return Reply( - content_type="text/event-stream", - chunks=cc.tool_use_stream(identity1, cc.FABLE, _THINKING, _SIGNATURE, calls, _USAGE), - ) - assert body == second_expected, _diff(second_expected, body) - return Reply( - content_type="text/event-stream", - chunks=cc.text_stream(identity2, cc.FABLE, "PROBE", {"input_tokens": 30, "output_tokens": 3}), - ) - - with wire_server(respond) as wire, gateway.scenario() as scenario: - model: Final = scenario.model(model=f"anthropic/{cc.FABLE}", api_base=wire.url, api_key=cc.ANTHROPIC_API_KEY) - headers: Final = cc.cli_headers(gateway.key, cc.FRONTIER_CLI_BETA) - response1: Final = gateway.request( - "POST", "/v1/messages", {**turn1, "model": model}, params={"beta": "true"}, headers=headers - ) - assert response1.status_code == 200, response1.text - events: Final = cc.sse_events(response1.text) - assert [ - ( - event, - data.get("delta", {}).get( - "type", data.get("content_block", {}).get("type", data.get("delta", {}).get("stop_reason")) - ), - ) - for event, data in events - ] == [ - ("message_start", None), - ("content_block_start", "thinking"), - ("content_block_delta", "thinking_delta"), - ("content_block_delta", "signature_delta"), - ("content_block_stop", None), - ("content_block_start", "tool_use"), - ("content_block_delta", "input_json_delta"), - ("content_block_delta", "input_json_delta"), - ("content_block_stop", None), - ("message_delta", "tool_use"), - ("message_stop", None), - ] - assert events[5][1]["content_block"]["id"] == "toolu_read_1" - assert events[5][1]["content_block"]["name"] == "Read" - partial: Final = events[6][1]["delta"]["partial_json"] + events[7][1]["delta"]["partial_json"] - assert partial == '{"file_path": "/tmp/cc_probe/hello.txt"}' - response2: Final = gateway.request( - "POST", "/v1/messages", {**turn2, "model": model}, params={"beta": "true"}, headers=headers - ) - assert response2.status_code == 200, response2.text - events2: Final = cc.sse_events(response2.text) - assert events2[2][1]["delta"] == {"type": "text_delta", "text": "PROBE"} - assert events2[4][1]["delta"]["stop_reason"] == "end_turn" - bodies: Final = tuple(cc.JSON_OBJECT.validate_json(request.body) for request in wire.drain()) - assert bodies == (first_expected, second_expected), bodies - rows: Final = eventually( - lambda: read_rows( - 'SELECT prompt_tokens FROM "LiteLLM_SpendLogs" WHERE request_id=%s', - (identity2,), - ), - lambda values: len(values) == 1, - seconds=70, - ) - assert rows[0]["prompt_tokens"] == 30 - - -def test_parallel_tool_results_reach_anthropic_in_client_order(gateway: Gateway) -> None: - turn1: Final = cc.frontier_request( - f"cache-bust-{uuid.uuid4().hex}", - "high", - 64000, - prompt_text="Read /tmp/cc_probe/hello.txt and /tmp/cc_probe/world.txt and reply with both words", - ) - turn2: Final = cc.tool_loop_turn2( - turn1, - ( - {"type": "thinking", "thinking": _THINKING, "signature": _SIGNATURE}, - { - "type": "tool_use", - "id": "toolu_read_1", - "name": "Read", - "input": {"file_path": "/tmp/cc_probe/hello.txt"}, - }, - { - "type": "tool_use", - "id": "toolu_read_2", - "name": "Read", - "input": {"file_path": "/tmp/cc_probe/world.txt"}, - }, - ), - (("toolu_read_2", "1\tPROBE2\n2\t"), ("toolu_read_1", "1\tPROBE\n2\t")), - ) - - def respond(request: Request) -> Reply: - body: Final = cc.JSON_OBJECT.validate_json(request.body) - expected: Final = {**turn2, "model": cc.FABLE} - assert body == expected, _diff(expected, body) - results: Final = [block for block in body["messages"][3]["content"] if block["type"] == "tool_result"] - assert [block["tool_use_id"] for block in results] == ["toolu_read_2", "toolu_read_1"] - return Reply( - content_type="text/event-stream", - chunks=cc.text_stream(f"msg_mt_{uuid.uuid4().hex}", cc.FABLE, "PROBE PROBE2", _USAGE), - ) - - with wire_server(respond) as wire, gateway.scenario() as scenario: - model: Final = scenario.model(model=f"anthropic/{cc.FABLE}", api_base=wire.url, api_key=cc.ANTHROPIC_API_KEY) - response: Final = gateway.request( - "POST", - "/v1/messages", - {**turn2, "model": model}, - params={"beta": "true"}, - headers=cc.cli_headers(gateway.key, cc.FRONTIER_CLI_BETA), - ) - assert response.status_code == 200, response.text - assert len(wire.drain()) == 1 diff --git a/tests/integration/messages_endpoint/providers/anthropic/tools/test_anthropic_web_search_citations_wire.py b/tests/integration/messages_endpoint/providers/anthropic/tools/test_anthropic_web_search_citations_wire.py deleted file mode 100644 index 66852e6e21b..00000000000 --- a/tests/integration/messages_endpoint/providers/anthropic/tools/test_anthropic_web_search_citations_wire.py +++ /dev/null @@ -1,170 +0,0 @@ -import uuid -from typing import Final - -from integration._support import claude_code as cc -from integration._support.client import Gateway, eventually -from integration._support.database import read_rows -from integration._support.wire import Reply, Request, wire_server - -WEB_SEARCH_TOOL: Final = {"type": "web_search_20250305", "name": "web_search", "max_uses": 8} - - -def _web_search_stream(identity: str) -> tuple[bytes, ...]: - return ( - cc.sse_frame( - "message_start", - { - "type": "message_start", - "message": { - "id": identity, - "type": "message", - "role": "assistant", - "model": cc.FABLE, - "content": [], - "stop_reason": None, - "stop_sequence": None, - "usage": {"input_tokens": 20, "output_tokens": 1, "server_tool_use": {"web_search_requests": 1}}, - }, - }, - ), - cc.sse_frame( - "content_block_start", - { - "type": "content_block_start", - "index": 0, - "content_block": {"type": "server_tool_use", "id": "srvtoolu_1", "name": "web_search", "input": {}}, - }, - ), - cc.sse_frame( - "content_block_delta", - { - "type": "content_block_delta", - "index": 0, - "delta": {"type": "input_json_delta", "partial_json": '{"query": "current LiteLLM version"}'}, - }, - ), - cc.sse_frame("content_block_stop", {"type": "content_block_stop", "index": 0}), - cc.sse_frame( - "content_block_start", - { - "type": "content_block_start", - "index": 1, - "content_block": { - "type": "web_search_tool_result", - "tool_use_id": "srvtoolu_1", - "content": [ - { - "type": "web_search_result", - "title": "litellm releases", - "url": "https://example.com/litellm", - "page_age": None, - "encrypted_content": "enc_ws_1", - } - ], - }, - }, - ), - cc.sse_frame("content_block_stop", {"type": "content_block_stop", "index": 1}), - cc.sse_frame( - "content_block_start", - {"type": "content_block_start", "index": 2, "content_block": {"type": "text", "text": ""}}, - ), - cc.sse_frame( - "content_block_delta", - {"type": "content_block_delta", "index": 2, "delta": {"type": "text_delta", "text": "1.104.0"}}, - ), - cc.sse_frame( - "content_block_delta", - { - "type": "content_block_delta", - "index": 2, - "delta": { - "type": "citations_delta", - "citation": { - "type": "web_search_result_location", - "url": "https://example.com/litellm", - "title": "litellm releases", - "cited_text": "version 1.104.0", - "encrypted_index": "eidx_1", - }, - }, - }, - ), - cc.sse_frame("content_block_stop", {"type": "content_block_stop", "index": 2}), - cc.sse_frame( - "message_delta", - { - "type": "message_delta", - "delta": {"stop_reason": "end_turn", "stop_sequence": None}, - "usage": {"output_tokens": 15, "server_tool_use": {"web_search_requests": 1}}, - }, - ), - cc.sse_frame("message_stop", {"type": "message_stop"}), - ) - - -def test_web_search_tool_passthrough_and_cited_response(gateway: Gateway) -> None: - identity: Final = f"msg_ws_{uuid.uuid4().hex}" - base: Final = cc.frontier_request( - f"cache-bust-{uuid.uuid4().hex}", - "high", - 64000, - prompt_text="Use web search to find the current LiteLLM version and answer in one word", - ) - request_body: Final = { - **base, - "tools": [*base["tools"], WEB_SEARCH_TOOL], - } - - def respond(request: Request) -> Reply: - assert request.method == "POST" - assert request.target == "/v1/messages", request.target - body: Final = cc.JSON_OBJECT.validate_json(request.body) - expected: Final = {**request_body, "model": cc.FABLE} - assert body == expected, { - key: (expected.get(key), body.get(key)) - for key in expected.keys() | body.keys() - if expected.get(key) != body.get(key) - } - assert body["tools"][-1] == WEB_SEARCH_TOOL - assert len({tool["name"] for tool in body["tools"]}) == len(body["tools"]), body["tools"] - return Reply(content_type="text/event-stream", chunks=_web_search_stream(identity)) - - with wire_server(respond) as wire, gateway.scenario() as scenario: - model: Final = scenario.model( - model=f"anthropic/{cc.FABLE}", - api_base=wire.url, - api_key=cc.ANTHROPIC_API_KEY, - input_cost_per_token=1e-6, - output_cost_per_token=5e-6, - ) - response: Final = gateway.request( - "POST", - "/v1/messages", - {**request_body, "model": model}, - params={"beta": "true"}, - headers=cc.cli_headers(gateway.key, cc.FRONTIER_CLI_BETA), - ) - assert response.status_code == 200, response.text - events: Final = cc.sse_events(response.text) - started: Final = [ - (data["index"], data["content_block"]["type"]) for event, data in events if event == "content_block_start" - ] - assert started == [(0, "server_tool_use"), (1, "web_search_tool_result"), (2, "text")], started - citations: Final = [ - data["delta"] for event, data in events if data.get("delta", {}).get("type") == "citations_delta" - ] - assert len(citations) == 1 and citations[0]["citation"]["url"] == "https://example.com/litellm", citations - start_usage: Final = events[0][1]["message"]["usage"] - assert start_usage["server_tool_use"]["web_search_requests"] == 1, start_usage - assert len(wire.drain()) == 1 - rows: Final = eventually( - lambda: read_rows( - 'SELECT spend, prompt_tokens, completion_tokens FROM "LiteLLM_SpendLogs" WHERE request_id=%s', - (identity,), - ), - lambda values: len(values) == 1, - seconds=70, - ) - token_cost: Final = 20 * 1e-6 + 15 * 5e-6 - assert float(rows[0]["spend"]) >= token_cost, dict(rows[0]) diff --git a/tests/integration/messages_endpoint/providers/anthropic/usage/test_anthropic_long_context_beta_wire.py b/tests/integration/messages_endpoint/providers/anthropic/usage/test_anthropic_long_context_beta_wire.py deleted file mode 100644 index 7fe405ca439..00000000000 --- a/tests/integration/messages_endpoint/providers/anthropic/usage/test_anthropic_long_context_beta_wire.py +++ /dev/null @@ -1,53 +0,0 @@ -import uuid -from typing import Final - -import pytest -from integration._support import claude_code as cc -from integration._support.client import Gateway, eventually -from integration._support.database import read_rows -from integration._support.wire import Reply, Request, wire_server - -_BETA_1M: Final = f"{cc.FRONTIER_CLI_BETA.replace(',effort-2025-11-24', ',context-1m-2025-08-07,effort-2025-11-24')}" - - -def test_1m_context_beta_forwarded_and_tiered_prompt_priced_above_200k(gateway: Gateway) -> None: - identity: Final = f"msg_1m_{uuid.uuid4().hex}" - request_body: Final = cc.frontier_request(f"cache-bust-{uuid.uuid4().hex}", "high", 64000) - - def respond(request: Request) -> Reply: - assert request.method == "POST" - assert request.target == "/v1/messages", request.target - upstream_beta: Final = request.headers.get("anthropic-beta", "") - assert upstream_beta.split(",").count("context-1m-2025-08-07") == 1, upstream_beta - return Reply( - content_type="text/event-stream", - chunks=cc.text_stream(identity, cc.FABLE, "PONG", {"input_tokens": 250000, "output_tokens": 100}), - ) - - with wire_server(respond) as wire, gateway.scenario() as scenario: - model: Final = scenario.model( - model=f"anthropic/{cc.FABLE}", - api_base=wire.url, - api_key=cc.ANTHROPIC_API_KEY, - input_cost_per_token=1e-6, - input_cost_per_token_above_200k_tokens=2e-6, - output_cost_per_token=5e-6, - ) - response: Final = gateway.request( - "POST", - "/v1/messages", - {**request_body, "model": model}, - params={"beta": "true"}, - headers=cc.cli_headers(gateway.key, _BETA_1M), - ) - assert response.status_code == 200, response.text - assert len(wire.drain()) == 1 - rows: Final = eventually( - lambda: read_rows( - 'SELECT spend, prompt_tokens, completion_tokens FROM "LiteLLM_SpendLogs" WHERE request_id=%s', - (identity,), - ), - lambda values: len(values) == 1, - seconds=70, - ) - assert float(rows[0]["spend"]) == pytest.approx(250000 * 2e-6 + 100 * 5e-6), dict(rows[0]) diff --git a/tests/integration/routing/test_overloaded_deployment_fallback.py b/tests/integration/routing/test_overloaded_deployment_fallback.py deleted file mode 100644 index 7f5e669472d..00000000000 --- a/tests/integration/routing/test_overloaded_deployment_fallback.py +++ /dev/null @@ -1,93 +0,0 @@ -import json -import uuid -from pathlib import Path -from typing import Final - -import yaml -from integration._support import claude_code as cc -from integration._support.client import Gateway, eventually -from integration._support.database import read_rows -from integration._support.process import owned_proxy -from integration._support.wire import Reply, Request, wire_server - - -def _error_529() -> Reply: - return Reply( - status=529, - body=json.dumps({"type": "error", "error": {"type": "overloaded_error", "message": "Overloaded"}}).encode(), - ) - - -def test_overloaded_primary_falls_back_to_second_deployment(gateway: Gateway, tmp_path: Path) -> None: - identity: Final = f"msg_fb_{uuid.uuid4().hex}" - request_body: Final = cc.claude_code_request(f"cache-bust-{uuid.uuid4().hex}") - - def respond_fallback(request: Request) -> Reply: - return Reply( - content_type="text/event-stream", - chunks=cc.text_stream(identity, cc.SONNET, "PONG", {"input_tokens": 12, "output_tokens": 4}), - ) - - with ( - wire_server(lambda request: _error_529()) as primary, - wire_server(respond_fallback) as fallback, - ): - config: Final = { - **yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text()), - "model_list": [ - { - "model_name": "cc-primary", - "litellm_params": { - "model": f"anthropic/{cc.SONNET}", - "api_key": cc.ANTHROPIC_API_KEY, - "api_base": primary.url, - "model_info": {"id": "primary-cc"}, - }, - }, - { - "model_name": "cc-fallback-group", - "litellm_params": { - "model": f"anthropic/{cc.SONNET}", - "api_key": cc.ANTHROPIC_API_KEY, - "api_base": fallback.url, - "model_info": {"id": "fallback-cc"}, - }, - }, - ], - "router_settings": { - "num_retries": 0, - "disable_cooldowns": True, - "fallbacks": [{"cc-primary": ["cc-fallback-group"]}], - }, - } - path: Final = tmp_path / "fallbacks.yaml" - path.write_text(yaml.safe_dump(config)) - with owned_proxy( - gateway, tmp_path, {"REDIS_HOST": "127.0.0.1", "REDIS_PORT": "6379"}, config=path - ) as candidate: - with candidate.client.stream( - "POST", - "/v1/messages", - params={"beta": "true"}, - json={**request_body, "model": "cc-primary"}, - headers={**cc.cli_headers(candidate.key), "authorization": f"Bearer {candidate.key}"}, - ) as response: - assert response.status_code == 200, response.status_code - body: Final = "".join(response.iter_text()) - assert "message_stop" in body, body - assert "PONG" in body, body - deployments: Final = candidate.get("/model/info")["data"] - fallback_id: Final = next( - entry["model_info"]["id"] - for entry in deployments - if entry["litellm_params"]["api_base"] == fallback.url - ) - assert response.headers.get("x-litellm-model-id") == fallback_id, dict(response.headers) - assert len(primary.drain()) == 1 - assert len(fallback.drain()) == 1 - rows: Final = eventually( - lambda: read_rows('SELECT request_id FROM "LiteLLM_SpendLogs" WHERE request_id=%s', (identity,)), - lambda values: len(values) == 1, - seconds=70, - ) - assert len(rows) == 1 diff --git a/tests/integration/streaming/test_client_disconnect_still_bills.py b/tests/integration/streaming/test_client_disconnect_still_bills.py deleted file mode 100644 index 19a52f274a1..00000000000 --- a/tests/integration/streaming/test_client_disconnect_still_bills.py +++ /dev/null @@ -1,85 +0,0 @@ -import uuid -from typing import Final - -from integration._support import claude_code as cc -from integration._support.client import Gateway, eventually -from integration._support.database import read_rows -from integration._support.wire import Reply, Request, wire_server - - -def test_client_disconnect_mid_stream_still_bills_the_message(gateway: Gateway) -> None: - identity: Final = f"msg_dc_{uuid.uuid4().hex}" - request_body: Final = cc.claude_code_request(f"cache-bust-{uuid.uuid4().hex}") - - def respond(request: Request) -> Reply: - return Reply( - content_type="text/event-stream", - chunks=( - cc.sse_frame( - "message_start", - { - "type": "message_start", - "message": { - "id": identity, - "type": "message", - "role": "assistant", - "model": cc.SONNET, - "content": [], - "stop_reason": None, - "stop_sequence": None, - "usage": {"input_tokens": 12, "output_tokens": 1}, - }, - }, - ), - cc.sse_frame( - "content_block_start", - {"type": "content_block_start", "index": 0, "content_block": {"type": "text", "text": ""}}, - ), - cc.sse_frame( - "content_block_delta", - {"type": "content_block_delta", "index": 0, "delta": {"type": "text_delta", "text": "PONG"}}, - ), - cc.sse_frame("content_block_stop", {"type": "content_block_stop", "index": 0}), - cc.sse_frame( - "message_delta", - { - "type": "message_delta", - "delta": {"stop_reason": "end_turn", "stop_sequence": None}, - "usage": {"output_tokens": 4}, - }, - ), - cc.sse_frame("message_stop", {"type": "message_stop"}), - ), - pause_between_chunks=3.0, - ) - - with wire_server(respond) as wire, gateway.scenario() as scenario: - model: Final = scenario.model(model=f"anthropic/{cc.SONNET}", api_base=wire.url, api_key=cc.ANTHROPIC_API_KEY) - with gateway.client.stream( - "POST", - "/v1/messages", - params={"beta": "true"}, - json={**request_body, "model": model}, - headers={ - **cc.cli_headers(gateway.key), - "authorization": f"Bearer {gateway.key}", - }, - ) as response: - assert response.status_code == 200, response.status_code - first: Final = next(response.iter_text()) - assert "message_start" in first, first - assert len(wire.drain()) == 1 - rows: Final = eventually( - lambda: read_rows( - 'SELECT prompt_tokens FROM "LiteLLM_SpendLogs" WHERE request_id=%s', - (identity,), - ), - lambda values: len(values) == 1, - seconds=70, - return_last_on_timeout=True, - ) - assert rows and rows[0]["prompt_tokens"] == 12, rows - assert wire.disconnected.empty(), ( - "closing the client stream must not abort the upstream call before it finishes; " - f"wire recorded a disconnect on {wire.disconnected.get_nowait()}" - )