mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-10 03:28:53 +00:00
Add @step labels to the HttpTransport methods, the poll and wait helpers and the boot helpers that did real IO without recording a step, so a test that reaches the proxy through them no longer reports an empty or gappy step timeline in the JUnit report.
487 lines
15 KiB
Python
487 lines
15 KiB
Python
"""Transport: the typed request primitives clients use, behind a Protocol.
|
|
|
|
`Transport` is what each client depends on (composition + DI); `HttpTransport` is
|
|
the concrete frozen-slots dataclass that fulfils it via the e2e_http wrapper. No
|
|
client touches requests.* or builds raw dicts; they pass pydantic models here.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from dataclasses import dataclass, field
|
|
from typing import Protocol
|
|
|
|
import e2e_http
|
|
from e2e_http import (
|
|
URL,
|
|
AbandonedRequest,
|
|
AuthHeaders,
|
|
BinaryStream,
|
|
NetworkError,
|
|
ProbeResult,
|
|
Result,
|
|
StreamHead,
|
|
StreamingResponse,
|
|
)
|
|
from e2e_metadata import step
|
|
from pydantic import BaseModel
|
|
|
|
|
|
class Transport(Protocol):
|
|
def post[R: BaseModel](
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
json: BaseModel,
|
|
response_type: type[R],
|
|
timeout: float | None = None,
|
|
) -> Result[R]: ...
|
|
|
|
def stream(self, path: str, *, headers: BaseModel, json: BaseModel) -> StreamingResponse: ...
|
|
|
|
def open_stream(self, path: str, *, headers: BaseModel, json: BaseModel) -> StreamHead | NetworkError: ...
|
|
|
|
def stream_binary(
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
json: BaseModel,
|
|
chunk_size: int = 8192,
|
|
) -> BinaryStream: ...
|
|
|
|
def send(
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
json: BaseModel,
|
|
params: BaseModel | None = None,
|
|
stream: bool = False,
|
|
) -> StreamingResponse: ...
|
|
|
|
def abandon(
|
|
self, path: str, *, headers: BaseModel, json: BaseModel, after: float
|
|
) -> AbandonedRequest | StreamingResponse: ...
|
|
|
|
def get[R: BaseModel](
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
params: BaseModel,
|
|
response_type: type[R],
|
|
timeout: float | None = None,
|
|
) -> Result[R]: ...
|
|
|
|
def delete[R: BaseModel](
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
json: BaseModel,
|
|
response_type: type[R],
|
|
params: BaseModel | None = None,
|
|
) -> Result[R]: ...
|
|
|
|
def patch[R: BaseModel](
|
|
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
|
|
) -> Result[R]: ...
|
|
|
|
def put[R: BaseModel](
|
|
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
|
|
) -> Result[R]: ...
|
|
|
|
def probe(self, path: str, *, params: BaseModel, headers: BaseModel | None = None) -> ProbeResult: ...
|
|
|
|
def upload[R: BaseModel](
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
form: BaseModel,
|
|
filename: str,
|
|
content: bytes,
|
|
file_content_type: str = "application/jsonl",
|
|
file_field: str = "file",
|
|
params: BaseModel | None = None,
|
|
response_type: type[R],
|
|
timeout: float | None = None,
|
|
) -> Result[R]: ...
|
|
|
|
def download(self, path: str, *, headers: BaseModel) -> StreamingResponse: ...
|
|
|
|
def bearer(self, key: str) -> AuthHeaders: ...
|
|
|
|
@property
|
|
def master(self) -> AuthHeaders: ...
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class HttpTransport:
|
|
base_url: str
|
|
master_key: str = field(repr=False)
|
|
request_timeout: float = 60.0
|
|
|
|
def _url(self, path: str) -> URL:
|
|
return URL(f"{self.base_url.rstrip('/')}{path}")
|
|
|
|
def bearer(self, key: str) -> AuthHeaders:
|
|
return AuthHeaders(authorization=f"Bearer {key}")
|
|
|
|
@property
|
|
def master(self) -> AuthHeaders:
|
|
return self.bearer(self.master_key)
|
|
|
|
@step("POST {path}")
|
|
def post[R: BaseModel](
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
json: BaseModel,
|
|
response_type: type[R],
|
|
timeout: float | None = None,
|
|
) -> Result[R]:
|
|
"""`timeout` overrides the transport-wide request_timeout for this call, for
|
|
provider operations that legitimately outlive it (image edits, OCR)."""
|
|
return e2e_http.post(
|
|
self._url(path),
|
|
headers=headers,
|
|
json=json,
|
|
response_type=response_type,
|
|
timeout=self.request_timeout if timeout is None else timeout,
|
|
)
|
|
|
|
@step("GET {path}")
|
|
def get[R: BaseModel](
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
params: BaseModel,
|
|
response_type: type[R],
|
|
timeout: float | None = None,
|
|
) -> Result[R]:
|
|
"""`timeout` overrides the transport-wide request_timeout for this call, for
|
|
pollers whose own deadline is shorter than it."""
|
|
return e2e_http.get(
|
|
self._url(path),
|
|
headers=headers,
|
|
params=params,
|
|
response_type=response_type,
|
|
timeout=self.request_timeout if timeout is None else timeout,
|
|
)
|
|
|
|
@step("DELETE {path}")
|
|
def delete[R: BaseModel](
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
json: BaseModel,
|
|
response_type: type[R],
|
|
params: BaseModel | None = None,
|
|
) -> Result[R]:
|
|
return e2e_http.delete(
|
|
self._url(path),
|
|
headers=headers,
|
|
json=json,
|
|
params=params,
|
|
response_type=response_type,
|
|
timeout=self.request_timeout,
|
|
)
|
|
|
|
@step("PATCH {path}")
|
|
def patch[R: BaseModel](
|
|
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
|
|
) -> Result[R]:
|
|
return e2e_http.patch(
|
|
self._url(path),
|
|
headers=headers,
|
|
json=json,
|
|
response_type=response_type,
|
|
timeout=self.request_timeout,
|
|
)
|
|
|
|
@step("PUT {path}")
|
|
def put[R: BaseModel](self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]) -> Result[R]:
|
|
return e2e_http.put(
|
|
self._url(path),
|
|
headers=headers,
|
|
json=json,
|
|
response_type=response_type,
|
|
timeout=self.request_timeout,
|
|
)
|
|
|
|
@step("Stream a POST to {path}")
|
|
def stream(self, path: str, *, headers: BaseModel, json: BaseModel) -> StreamingResponse:
|
|
return e2e_http.stream(self._url(path), headers=headers, json=json, timeout=self.request_timeout)
|
|
|
|
@step("Open a stream to {path}")
|
|
def open_stream(self, path: str, *, headers: BaseModel, json: BaseModel) -> StreamHead | NetworkError:
|
|
return e2e_http.open_stream(self._url(path), headers=headers, json=json, timeout=self.request_timeout)
|
|
|
|
@step("Stream binary from {path}")
|
|
def stream_binary(
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
json: BaseModel,
|
|
chunk_size: int = 8192,
|
|
) -> BinaryStream:
|
|
return e2e_http.stream_binary(
|
|
self._url(path),
|
|
headers=headers,
|
|
json=json,
|
|
chunk_size=chunk_size,
|
|
timeout=self.request_timeout,
|
|
)
|
|
|
|
@step("Send a request to {path}")
|
|
def send(
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
json: BaseModel,
|
|
params: BaseModel | None = None,
|
|
stream: bool = False,
|
|
) -> StreamingResponse:
|
|
return e2e_http.send(
|
|
self._url(path),
|
|
headers=headers,
|
|
json=json,
|
|
params=params,
|
|
stream=stream,
|
|
timeout=self.request_timeout,
|
|
)
|
|
|
|
@step("Abandon the request to {path} after {after}s")
|
|
def abandon(
|
|
self, path: str, *, headers: BaseModel, json: BaseModel, after: float
|
|
) -> AbandonedRequest | StreamingResponse:
|
|
return e2e_http.abandon(self._url(path), headers=headers, json=json, after=after)
|
|
|
|
@step("Probe {path}")
|
|
def probe(self, path: str, *, params: BaseModel, headers: BaseModel | None = None) -> ProbeResult:
|
|
return e2e_http.probe(
|
|
self._url(path),
|
|
headers=self.master if headers is None else headers,
|
|
params=params,
|
|
timeout=self.request_timeout,
|
|
)
|
|
|
|
@step("Upload {filename} to {path}")
|
|
def upload[R: BaseModel](
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
form: BaseModel,
|
|
filename: str,
|
|
content: bytes,
|
|
file_content_type: str = "application/jsonl",
|
|
file_field: str = "file",
|
|
params: BaseModel | None = None,
|
|
response_type: type[R],
|
|
timeout: float | None = None,
|
|
) -> Result[R]:
|
|
return e2e_http.upload(
|
|
self._url(path),
|
|
headers=headers,
|
|
form=form,
|
|
filename=filename,
|
|
content=content,
|
|
file_content_type=file_content_type,
|
|
file_field=file_field,
|
|
params=params,
|
|
response_type=response_type,
|
|
timeout=self.request_timeout if timeout is None else timeout,
|
|
)
|
|
|
|
@step("Download {path}")
|
|
def download(self, path: str, *, headers: BaseModel) -> StreamingResponse:
|
|
return e2e_http.download(self._url(path), headers=headers, timeout=self.request_timeout)
|
|
|
|
|
|
# Top-level management/admin route groups. In a split deployment these are served
|
|
# by the control plane (a different service from the LLM data plane). LLM routes
|
|
# (/chat, /embeddings, and native passthrough like /gemini, /anthropic) are NOT
|
|
# here and fall through to the data plane. Matched as path prefixes.
|
|
CONTROL_PLANE_PREFIXES: tuple[str, ...] = (
|
|
"/key",
|
|
"/user",
|
|
"/team",
|
|
"/organization",
|
|
"/project",
|
|
"/customer",
|
|
"/end_user",
|
|
"/tag",
|
|
"/budget",
|
|
"/model/",
|
|
"/access_group",
|
|
"/spend",
|
|
"/global",
|
|
"/config",
|
|
"/guardrails",
|
|
"/credentials",
|
|
"/router/settings",
|
|
"/audit",
|
|
"/public",
|
|
"/v2/login",
|
|
"/openapi.json",
|
|
)
|
|
|
|
|
|
def is_control_plane_path(path: str) -> bool:
|
|
"""True if `path` is a management/admin route (served by the control plane in a
|
|
split deployment), false for LLM data-plane routes."""
|
|
return path.startswith(CONTROL_PLANE_PREFIXES)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class SplitTransport:
|
|
"""A Transport that dispatches each call by path to one of two backends: the
|
|
management/admin control plane or the LLM data plane.
|
|
|
|
Litellm can run as a split control-plane/data-plane deployment where the two
|
|
surfaces live on different services. Clients here stay plane-agnostic — they
|
|
keep calling ``transport.post("/budget/new", ...)`` or
|
|
``transport.send("/chat/completions", ...)`` — and routing happens in one place
|
|
by path (see ``CONTROL_PLANE_PREFIXES``). When ``control`` and ``data`` share a
|
|
base URL (the monolithic default), routing is a no-op. ``bearer``/``master``
|
|
are plane-agnostic (same master key both planes), so they come from ``data``.
|
|
"""
|
|
|
|
data: HttpTransport
|
|
control: HttpTransport
|
|
|
|
def _route(self, path: str) -> HttpTransport:
|
|
return self.control if is_control_plane_path(path) else self.data
|
|
|
|
def bearer(self, key: str) -> AuthHeaders:
|
|
return self.data.bearer(key)
|
|
|
|
@property
|
|
def master(self) -> AuthHeaders:
|
|
return self.data.master
|
|
|
|
def post[R: BaseModel](
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
json: BaseModel,
|
|
response_type: type[R],
|
|
timeout: float | None = None,
|
|
) -> Result[R]:
|
|
return self._route(path).post(path, headers=headers, json=json, response_type=response_type, timeout=timeout)
|
|
|
|
def get[R: BaseModel](
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
params: BaseModel,
|
|
response_type: type[R],
|
|
timeout: float | None = None,
|
|
) -> Result[R]:
|
|
return self._route(path).get(
|
|
path,
|
|
headers=headers,
|
|
params=params,
|
|
response_type=response_type,
|
|
timeout=timeout,
|
|
)
|
|
|
|
def delete[R: BaseModel](
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
json: BaseModel,
|
|
response_type: type[R],
|
|
params: BaseModel | None = None,
|
|
) -> Result[R]:
|
|
return self._route(path).delete(
|
|
path,
|
|
headers=headers,
|
|
json=json,
|
|
response_type=response_type,
|
|
params=params,
|
|
)
|
|
|
|
def patch[R: BaseModel](
|
|
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
|
|
) -> Result[R]:
|
|
return self._route(path).patch(path, headers=headers, json=json, response_type=response_type)
|
|
|
|
def put[R: BaseModel](self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]) -> Result[R]:
|
|
return self._route(path).put(path, headers=headers, json=json, response_type=response_type)
|
|
|
|
def stream(self, path: str, *, headers: BaseModel, json: BaseModel) -> StreamingResponse:
|
|
return self._route(path).stream(path, headers=headers, json=json)
|
|
|
|
def open_stream(self, path: str, *, headers: BaseModel, json: BaseModel) -> StreamHead | NetworkError:
|
|
return self._route(path).open_stream(path, headers=headers, json=json)
|
|
|
|
def stream_binary(
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
json: BaseModel,
|
|
chunk_size: int = 8192,
|
|
) -> BinaryStream:
|
|
return self._route(path).stream_binary(path, headers=headers, json=json, chunk_size=chunk_size)
|
|
|
|
def send(
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
json: BaseModel,
|
|
params: BaseModel | None = None,
|
|
stream: bool = False,
|
|
) -> StreamingResponse:
|
|
return self._route(path).send(path, headers=headers, json=json, params=params, stream=stream)
|
|
|
|
def abandon(
|
|
self, path: str, *, headers: BaseModel, json: BaseModel, after: float
|
|
) -> AbandonedRequest | StreamingResponse:
|
|
return self._route(path).abandon(path, headers=headers, json=json, after=after)
|
|
|
|
def probe(self, path: str, *, params: BaseModel, headers: BaseModel | None = None) -> ProbeResult:
|
|
return self._route(path).probe(path, params=params, headers=headers)
|
|
|
|
def upload[R: BaseModel](
|
|
self,
|
|
path: str,
|
|
*,
|
|
headers: BaseModel,
|
|
form: BaseModel,
|
|
filename: str,
|
|
content: bytes,
|
|
file_content_type: str = "application/jsonl",
|
|
file_field: str = "file",
|
|
params: BaseModel | None = None,
|
|
response_type: type[R],
|
|
timeout: float | None = None,
|
|
) -> Result[R]:
|
|
return self._route(path).upload(
|
|
path,
|
|
headers=headers,
|
|
form=form,
|
|
filename=filename,
|
|
content=content,
|
|
file_content_type=file_content_type,
|
|
file_field=file_field,
|
|
params=params,
|
|
response_type=response_type,
|
|
timeout=timeout,
|
|
)
|
|
|
|
def download(self, path: str, *, headers: BaseModel) -> StreamingResponse:
|
|
return self._route(path).download(path, headers=headers)
|