mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-10 22:41:41 +00:00
The lifecycle suite polled /model/info on every URL in PROXY_REPLICA_URLS. Those URLs are the stack's gateways, and gateway/routes/allowlist.py trims them to the LLM data-plane surface, so /model/info answers only on the backend and 404s on every replica. All five tests failed at their first read-back in CI while passing against a monolith, where one process serves both planes. The stored row has one answer behind it, so it is read through the shared transport, which routes control-plane paths to the backend. What every gateway must agree on is which models it serves, so the create and delete steps poll /v1/models per replica instead, a route the gateway does serve. read_back_everywhere now rejects a control-plane path outright rather than timing out on it. Two things surfaced behind that. /public/ was missing from the transport's control-plane prefixes, so model_cost_map() was routed to a gateway and 404'd, and the billing steps needed a data-plane wait: a PATCH lands on the backend and each gateway picks it up on its own config reload, measured here at 12-24s, so they now drive calls until the new rate reaches the spend row and let the deadline fail them. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01C1S92J8gSxxKVe1JBzxWBF
468 lines
13 KiB
Python
468 lines
13 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
|
|
from typing import Protocol
|
|
|
|
from pydantic import BaseModel
|
|
|
|
import e2e_http
|
|
from e2e_http import (
|
|
URL,
|
|
AuthHeaders,
|
|
BinaryStream,
|
|
ProbeResult,
|
|
Result,
|
|
StreamingResponse,
|
|
)
|
|
|
|
|
|
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 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 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) -> 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
|
|
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)
|
|
|
|
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,
|
|
)
|
|
|
|
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,
|
|
)
|
|
|
|
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,
|
|
)
|
|
|
|
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,
|
|
)
|
|
|
|
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,
|
|
)
|
|
|
|
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
|
|
)
|
|
|
|
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,
|
|
)
|
|
|
|
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,
|
|
)
|
|
|
|
def probe(self, path: str, *, params: BaseModel) -> ProbeResult:
|
|
return e2e_http.probe(
|
|
self._url(path),
|
|
headers=self.master,
|
|
params=params,
|
|
timeout=self.request_timeout,
|
|
)
|
|
|
|
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,
|
|
)
|
|
|
|
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",
|
|
"/customer",
|
|
"/end_user",
|
|
"/tag",
|
|
"/budget",
|
|
"/model/",
|
|
"/access_group",
|
|
"/spend",
|
|
"/global",
|
|
"/config",
|
|
"/guardrails",
|
|
"/openapi.json",
|
|
"/public/",
|
|
)
|
|
|
|
|
|
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 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 probe(self, path: str, *, params: BaseModel) -> ProbeResult:
|
|
return self._route(path).probe(path, params=params)
|
|
|
|
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)
|