mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-11 03:38:38 +00:00
* feat(prometheus): cap series per metric for every labeled metric Add prometheus_metrics_max_series_per_metric: per metric and per worker process, the first N label sets keep a series of their own. Counters and histograms record every later label set on one series whose labels are all "other", so totals stay exact, and gauges skip it. The cap holds with multiple workers because it never needs to remove a series. Add prometheus_metrics_ttl_seconds: a series idle for that long is removed and its slot is freed. The prometheus client cannot remove a series in multi-process mode, so the TTL is ignored there with a startup warning. Both settings are off by default. The end_user caps are unchanged. * fix(prometheus): share the series cap across workers of one proxy instance Workers writing to one PROMETHEUS_MULTIPROC_DIR now agree on which label sets get a series through an append-only admissions file per metric, so a merged scrape stays at the cap plus `other` instead of growing with every worker and every worker restart. The two fallback counters now pass their label names as a keyword so the cap and prometheus_exclude_labels apply to them, admission and child creation happen under one lock, the test fixture restores the shared registry, and the `other` label value lives in constants.py. * test(prometheus): check emitted labels instead of wrapper types, close the admission match The exclude-labels test now emits through the spend and provider budget metrics and checks the scrape keeps all their labels. The admission match arms end in assert_never so the match is exhaustive. * fix(prometheus): return the exhaustive-match fallback so every admission arm returns * fix(prometheus): pick the series tracker with isinstance so every path of _admits returns * fix(prometheus): skip an admissions line a worker could only write part of * fix(prometheus): frame each admissions record with newlines so a cut-off record cannot swallow the next A record a worker could only write part of used to merge with the next worker's record, and both were skipped for one request. Each record is now written between two newlines, so the fragment is a line of its own. The clock fixture in the series tests starts from a constant instead of reading the real clock * fix(prometheus): ignore a non-positive series cap or TTL with a warning instead of failing the logger A cap or TTL of 0 or less raised at logger init. The proxy logs that as a non-blocking error and keeps serving, so the result was a running proxy with no Prometheus metrics at all. The setting is now ignored with a startup warning naming it, the same rule the end_user cap already follows for a non-positive value * fix(prometheus): start the series cap over on a one-worker restart and audit it live A proxy with one worker and an operator-set PROMETHEUS_MULTIPROC_DIR now drops litellm's admission files at boot, so a restart frees every slot there the way it already does with several workers. A cap or TTL that is not a number greater than 0 (a bool, a non-numeric string, an empty value) is ignored with the startup warning instead of breaking the logger The integration cells drive the cap on every endpoint through the OpenAI and Anthropic SDKs and raw httpx, streaming and not, plus gauges, cache hits, failures, both workers of one instance, the TTL on one worker and its warning on two, ignored settings, excluded labels on the fallback counters, a null cap, /config/update, a concurrent burst scraped mid-flight, a provider outage, a killed worker, and restarts with one and two workers * fix(prometheus): wipe an operator-set multiprocess directory on a one-worker boot too * fix(prometheus): leave the multiprocess directory alone on a setup-only run A run with --skip_server_startup starts no worker, so it no longer creates or wipes PROMETHEUS_MULTIPROC_DIR. Wiping there deleted the samples of a proxy already running against the same directory * fix(prometheus): free the capped series slots when a gateway or backend container restarts The component image entrypoint starts uvicorn without the proxy CLI and wiped only the .db sample files at container start, so the admitted-series files of the previous container survived an in-place restart. Every label set seen after the restart was then counted on `other` once the previous container had filled the cap * test(prometheus): cover a setup-only run and a gateway image restart under the cap Two integration cells from the audit: a `--skip_server_startup` run pointed at a live two-worker proxy's operator directory leaves its samples alone, and the gateway image (`docker/component_entrypoint.sh` running `python -m gateway.launch`) restarted on a kept PROMETHEUS_MULTIPROC_DIR starts the cap over. The burst cells now wait for every counter they assert on, since the request and failure counters of one call increment at different points of the logging callback * test(prometheus): prove the cap reaches the fallback counters in the X1 cell * fix(prometheus): ignore a cleanup interval that is not a number of at least 0 A string or negative prometheus_metrics_cleanup_interval_seconds reached the series tracker unvalidated, so the first labeled emit with a TTL on raised TypeError inside the callback and recorded no series. The interval is now validated the way the cap and the TTL are: an invalid value is ignored with a warning and the default 60 seconds applies. The I2 integration cell drives a string interval through a live proxy and reads the warning from its log * test(prometheus): give the restart cells the boot budget of their siblings C4 and C5 boot two proxies each and hit the file's 240 s budget on a loaded box; C3 and D1 already carry 420 s --------- Co-authored-by: mateo-berri <277851410+mateo-berri@users.noreply.github.com>
408 lines
13 KiB
Python
408 lines
13 KiB
Python
import errno
|
|
import os
|
|
import signal
|
|
import socket
|
|
import subprocess
|
|
import sys
|
|
import time
|
|
import uuid
|
|
from collections.abc import Generator, Iterator, Mapping
|
|
from contextlib import contextmanager
|
|
from dataclasses import dataclass
|
|
from pathlib import Path
|
|
from types import MappingProxyType
|
|
from typing import Final
|
|
|
|
import httpx
|
|
import psutil
|
|
from integration._support.client import GATEWAY_LIMITS, Gateway
|
|
|
|
DB_PUSH: Final = ("--use_prisma_db_push",)
|
|
MIGRATE_DEPLOY: Final = ()
|
|
LEGACY_MIGRATE_DEPLOY: Final = ("--use_legacy_migration_resolver",)
|
|
|
|
|
|
def proxy_database_environment() -> Mapping[str, str]:
|
|
writer: Final = os.environ.get("INTEGRATION_PROXY_DATABASE_URL", "")
|
|
reader: Final = os.environ.get("INTEGRATION_PROXY_READ_REPLICA_URL", "")
|
|
return MappingProxyType(
|
|
{
|
|
**({"DATABASE_URL": writer} if writer else {}),
|
|
**({"DATABASE_URL_READ_REPLICA": reader} if reader else {}),
|
|
}
|
|
)
|
|
|
|
|
|
def in_group(process: psutil.Process, group: int) -> bool:
|
|
try:
|
|
return os.getpgid(process.pid) == group
|
|
except ProcessLookupError:
|
|
return False
|
|
|
|
|
|
def group_members(group: int) -> tuple[psutil.Process, ...]:
|
|
return tuple(process for process in psutil.process_iter() if in_group(process, group))
|
|
|
|
|
|
def signal_group(group: int, action: int) -> None:
|
|
try:
|
|
os.killpg(group, action)
|
|
except ProcessLookupError:
|
|
pass
|
|
|
|
|
|
def graceful_stop_seconds() -> float:
|
|
return max(30.0, float(os.environ.get("INTEGRATION_PROXY_READY_SECONDS", "70")))
|
|
|
|
|
|
def stop_root_process(process: subprocess.Popen[bytes]) -> bool:
|
|
if process.poll() is not None:
|
|
return True
|
|
process.terminate()
|
|
try:
|
|
process.wait(timeout=graceful_stop_seconds())
|
|
except subprocess.TimeoutExpired:
|
|
return False
|
|
return True
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class OwnedProxy:
|
|
gateway: Gateway
|
|
process: subprocess.Popen[bytes]
|
|
log: Path
|
|
|
|
|
|
@contextmanager
|
|
def owned_proxy(
|
|
gateway: Gateway,
|
|
directory: Path,
|
|
overrides: Mapping[str, str],
|
|
*,
|
|
config: Path | None = None,
|
|
remove_environment: tuple[str, ...] = (),
|
|
workers: int = 1,
|
|
database_setup: tuple[str, ...] = DB_PUSH,
|
|
) -> Iterator[Gateway]:
|
|
with owned_proxy_process(
|
|
gateway,
|
|
directory,
|
|
overrides,
|
|
config=config,
|
|
remove_environment=remove_environment,
|
|
workers=workers,
|
|
database_setup=database_setup,
|
|
) as owned:
|
|
yield owned.gateway
|
|
|
|
|
|
def _stop(process: subprocess.Popen[bytes]) -> None:
|
|
root_stopped: Final = stop_root_process(process)
|
|
residual: Final = group_members(process.pid)
|
|
if residual:
|
|
signal_group(process.pid, signal.SIGTERM)
|
|
psutil.wait_procs(residual, timeout=5)
|
|
remaining: Final = group_members(process.pid)
|
|
if remaining:
|
|
signal_group(process.pid, signal.SIGKILL)
|
|
psutil.wait_procs(remaining, timeout=3)
|
|
process.wait(timeout=3)
|
|
survivors: Final = group_members(process.pid)
|
|
assert not survivors, "Owned proxy child survived cleanup"
|
|
assert root_stopped and not remaining, "Owned proxy required forced cleanup"
|
|
|
|
|
|
_PORT_ATTEMPTS: Final = 3
|
|
_BIND_COLLISION: Final = os.strerror(errno.EADDRINUSE)
|
|
|
|
|
|
def _free_port() -> int:
|
|
with socket.socket() as reserve:
|
|
reserve.bind(("127.0.0.1", 0))
|
|
return reserve.getsockname()[1]
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class _Launch:
|
|
process: subprocess.Popen[bytes]
|
|
port: int
|
|
log: Path
|
|
|
|
|
|
def _launch(command: tuple[str, ...], root: Path, environment: Mapping[str, str], output: Path) -> _Launch:
|
|
port: Final = _free_port()
|
|
log_path: Final = output / f"owned-proxy-{uuid.uuid4().hex}.log"
|
|
with log_path.open("w") as log:
|
|
process: Final = subprocess.Popen(
|
|
[*command, "--port", str(port)],
|
|
cwd=root,
|
|
env=environment,
|
|
stdout=log,
|
|
stderr=subprocess.STDOUT,
|
|
start_new_session=True,
|
|
)
|
|
return _Launch(process, port, log_path)
|
|
|
|
|
|
def _lost_port_race(exit_code: int | None, log: Path) -> bool:
|
|
return exit_code is not None and _BIND_COLLISION in log.read_text()
|
|
|
|
|
|
def _wait_until_ready(launch: _Launch) -> None:
|
|
with httpx.Client(base_url=f"http://127.0.0.1:{launch.port}", timeout=15, trust_env=False) as client:
|
|
deadline: Final = time.monotonic() + float(os.environ.get("INTEGRATION_PROXY_READY_SECONDS", "70"))
|
|
while launch.process.poll() is None:
|
|
try:
|
|
if client.get("/health/readiness", timeout=2).status_code == 200:
|
|
return
|
|
except httpx.TransportError:
|
|
pass
|
|
assert time.monotonic() < deadline, "Owned proxy readiness deadline exceeded"
|
|
time.sleep(0.1)
|
|
|
|
|
|
def _launch_until_bound(
|
|
command: tuple[str, ...], root: Path, environment: Mapping[str, str], output: Path, attempts: int
|
|
) -> _Launch:
|
|
launch: Final = _launch(command, root, environment, output)
|
|
try:
|
|
_wait_until_ready(launch)
|
|
exit_code: Final = launch.process.poll()
|
|
assert exit_code is None or (attempts > 1 and _lost_port_race(exit_code, launch.log)), (
|
|
"Owned proxy exited before readiness"
|
|
)
|
|
except BaseException:
|
|
_stop(launch.process)
|
|
raise
|
|
if exit_code is None:
|
|
return launch
|
|
_stop(launch.process)
|
|
return _launch_until_bound(command, root, environment, output, attempts - 1)
|
|
|
|
|
|
def _proxy_root() -> Path:
|
|
return Path(os.environ.get("INTEGRATION_PROXY_ROOT") or Path(__file__).resolve().parents[3])
|
|
|
|
|
|
def _proxy_environment(
|
|
gateway: Gateway, overrides: Mapping[str, str], remove_environment: tuple[str, ...]
|
|
) -> Mapping[str, str]:
|
|
return MappingProxyType(
|
|
{
|
|
**{
|
|
name: value
|
|
for name, value in {**os.environ, **proxy_database_environment()}.items()
|
|
if name not in remove_environment
|
|
},
|
|
"LITELLM_MASTER_KEY": gateway.key,
|
|
"LITELLM_SALT_KEY": os.environ.get("LITELLM_SALT_KEY", "sk-integration-salt"),
|
|
"STORE_MODEL_IN_DB": "True",
|
|
**overrides,
|
|
}
|
|
)
|
|
|
|
|
|
def setup_only_proxy_run(
|
|
gateway: Gateway, overrides: Mapping[str, str], *, config: Path, workers: int
|
|
) -> subprocess.CompletedProcess[str]:
|
|
"""The proxy CLI's `--skip_server_startup` pass (the image's setup step), run to completion with the
|
|
environment an owned proxy gets."""
|
|
return subprocess.run( # test-quality-ok: the checkout at the working directory is the proxy under test
|
|
(
|
|
sys.executable,
|
|
"-m",
|
|
"integration._support.proxy",
|
|
"--config",
|
|
str(config),
|
|
"--num_workers",
|
|
str(workers),
|
|
*DB_PUSH,
|
|
"--skip_server_startup",
|
|
),
|
|
cwd=_proxy_root(),
|
|
env=dict(_proxy_environment(gateway, overrides, ())),
|
|
capture_output=True,
|
|
text=True,
|
|
timeout=300,
|
|
check=False,
|
|
)
|
|
|
|
|
|
@contextmanager
|
|
def owned_proxy_process(
|
|
gateway: Gateway,
|
|
directory: Path,
|
|
overrides: Mapping[str, str],
|
|
*,
|
|
config: Path | None = None,
|
|
remove_environment: tuple[str, ...] = (),
|
|
workers: int = 1,
|
|
database_setup: tuple[str, ...] = DB_PUSH,
|
|
extra_arguments: tuple[str, ...] = (),
|
|
) -> Iterator[OwnedProxy]:
|
|
root: Final = _proxy_root()
|
|
environment: Final = _proxy_environment(gateway, overrides, remove_environment)
|
|
output: Final = Path(os.environ.get("INTEGRATION_RESULTS_DIR", str(directory)))
|
|
output.mkdir(parents=True, exist_ok=True)
|
|
command: Final = (
|
|
sys.executable,
|
|
"-m",
|
|
"integration._support.proxy",
|
|
"--config",
|
|
str(config or "tests/integration/proxy_config.yaml"),
|
|
"--host",
|
|
"127.0.0.1",
|
|
"--num_workers",
|
|
str(workers),
|
|
*database_setup,
|
|
*extra_arguments,
|
|
)
|
|
launch: Final = _launch_until_bound(command, root, environment, output, _PORT_ATTEMPTS)
|
|
process: Final = launch.process
|
|
try:
|
|
with httpx.Client(
|
|
base_url=f"http://127.0.0.1:{launch.port}", timeout=15, trust_env=False, limits=GATEWAY_LIMITS
|
|
) as client:
|
|
yield OwnedProxy(Gateway(client, gateway.key, gateway.upstream_url), process, launch.log)
|
|
finally:
|
|
_stop(process)
|
|
|
|
|
|
@contextmanager
|
|
def owned_gateway_image(
|
|
gateway: Gateway, directory: Path, overrides: Mapping[str, str], *, config: Path, workers: int
|
|
) -> Iterator[OwnedProxy]:
|
|
"""The componentized gateway started the way its image starts it: `docker/component_entrypoint.sh` running
|
|
`python -m gateway.launch`, with the config handed over as `CONFIG_FILE_PATH`. It serves the data plane only,
|
|
so keys come from a proxy that shares its database."""
|
|
root: Final = _proxy_root()
|
|
environment: Final = _proxy_environment(gateway, {**overrides, "CONFIG_FILE_PATH": str(config)}, ())
|
|
output: Final = Path(os.environ.get("INTEGRATION_RESULTS_DIR", str(directory)))
|
|
output.mkdir(parents=True, exist_ok=True)
|
|
command: Final = (
|
|
str(root / "docker" / "component_entrypoint.sh"),
|
|
sys.executable,
|
|
"-m",
|
|
"gateway.launch",
|
|
"--workers",
|
|
str(workers),
|
|
"--host",
|
|
"127.0.0.1",
|
|
)
|
|
launch: Final = _launch_until_bound(command, root, environment, output, _PORT_ATTEMPTS)
|
|
try:
|
|
with httpx.Client(
|
|
base_url=f"http://127.0.0.1:{launch.port}", timeout=15, trust_env=False, limits=GATEWAY_LIMITS
|
|
) as client:
|
|
yield OwnedProxy(Gateway(client, gateway.key, gateway.upstream_url), launch.process, launch.log)
|
|
finally:
|
|
_stop(launch.process)
|
|
|
|
|
|
def _is_ready(client: httpx.Client) -> bool:
|
|
try:
|
|
return client.get("/health/readiness", timeout=2).status_code == 200
|
|
except httpx.TransportError:
|
|
return False
|
|
|
|
|
|
def refused_boot_log(
|
|
gateway: Gateway,
|
|
directory: Path,
|
|
overrides: Mapping[str, str],
|
|
*,
|
|
config: Path | None = None,
|
|
) -> str:
|
|
"""Start the proxy and return its log once it exits non-zero instead of becoming ready."""
|
|
root: Final = _proxy_root()
|
|
environment: Final = _proxy_environment(gateway, overrides, ())
|
|
output: Final = Path(os.environ.get("INTEGRATION_RESULTS_DIR", str(directory)))
|
|
output.mkdir(parents=True, exist_ok=True)
|
|
command: Final = (
|
|
sys.executable,
|
|
"-m",
|
|
"integration._support.proxy",
|
|
"--config",
|
|
str(config or "tests/integration/proxy_config.yaml"),
|
|
"--host",
|
|
"127.0.0.1",
|
|
"--num_workers",
|
|
"1",
|
|
*DB_PUSH,
|
|
)
|
|
launch: Final = _launch(command, root, environment, output)
|
|
try:
|
|
with httpx.Client(base_url=f"http://127.0.0.1:{launch.port}", timeout=15, trust_env=False) as client:
|
|
deadline: Final = time.monotonic() + 70
|
|
while launch.process.poll() is None:
|
|
assert not _is_ready(client), (
|
|
f"Proxy became ready instead of refusing to boot:\n{launch.log.read_text()}"
|
|
)
|
|
assert time.monotonic() < deadline, "Proxy neither exited nor became ready within the deadline"
|
|
time.sleep(0.1)
|
|
assert launch.process.returncode != 0, f"Proxy exited 0 instead of refusing to boot:\n{launch.log.read_text()}"
|
|
return launch.log.read_text()
|
|
finally:
|
|
_stop(launch.process)
|
|
|
|
|
|
_UPSTREAM_READY_SECONDS: Final = 60
|
|
|
|
|
|
class UpstreamSlot:
|
|
"""A scripted upstream a test module owns on a fixed port, so a cell can take it down and bring it back."""
|
|
|
|
__slots__ = ("directory", "port", "process", "root")
|
|
|
|
def __init__(self, directory: Path, port: int, root: Path) -> None:
|
|
self.directory = directory
|
|
self.port = port
|
|
self.root = root
|
|
self.process: subprocess.Popen[bytes] | None = None
|
|
|
|
@property
|
|
def url(self) -> str:
|
|
return f"http://127.0.0.1:{self.port}"
|
|
|
|
def start(self) -> None:
|
|
assert self.process is None, "Owned upstream is already running"
|
|
output: Final = Path(os.environ.get("INTEGRATION_RESULTS_DIR") or self.directory)
|
|
log_path: Final = output / f"owned-upstream-{self.port}-{uuid.uuid4().hex}.log"
|
|
with log_path.open("w") as log:
|
|
process: Final = subprocess.Popen(
|
|
[sys.executable, "-m", "integration._support.upstream", "--port", str(self.port)],
|
|
cwd=self.root,
|
|
env=dict(os.environ),
|
|
stdout=log,
|
|
stderr=subprocess.STDOUT,
|
|
start_new_session=True,
|
|
)
|
|
self.process = process
|
|
deadline: Final = time.monotonic() + _UPSTREAM_READY_SECONDS
|
|
while process.poll() is None:
|
|
try:
|
|
if httpx.get(f"{self.url}/health", timeout=2, trust_env=False).status_code == 200:
|
|
return
|
|
except httpx.TransportError:
|
|
pass
|
|
assert time.monotonic() < deadline, f"Owned upstream readiness deadline exceeded: {log_path}"
|
|
time.sleep(0.1)
|
|
raise AssertionError(f"Owned upstream exited before readiness: {log_path}")
|
|
|
|
def stop(self) -> None:
|
|
process: Final = self.process
|
|
assert process is not None, "Owned upstream is not running"
|
|
self.process = None
|
|
_stop(process)
|
|
|
|
|
|
@contextmanager
|
|
def owned_upstream(directory: Path) -> Generator[UpstreamSlot]:
|
|
root: Final = _proxy_root()
|
|
slot: Final = UpstreamSlot(directory, _free_port(), root)
|
|
slot.start()
|
|
try:
|
|
yield slot
|
|
finally:
|
|
if slot.process is not None:
|
|
slot.stop()
|