litellm/tests/integration/_support/process.py
devin-ai-integration[bot] 8b1990b4bc
feat(decisions): add unified /v1/decisions endpoint for Jev-compatible providers (#44236)
* feat(decisions): add unified /v1/decisions endpoint for Jev-compatible providers

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* fix(decisions): register typesafe as a provider so Jev deployments load

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* refactor(decisions): move provider endpoints under llms and validate proxy bodies

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* feat(decisions): add Cloudflare Clef and Strands Decider backends

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* fix(decisions): register decisions routes for managed agents and gateway

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* test(decisions): use raw regex for cloudflare missing account match

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* fix(decisions): avoid cast in Cloudflare response unwrapping

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* fix(decisions): default model, evaluation health probe, short Cloudflare names

The proxy validates only state and questions, so a request without a
model falls through to the configured default model like every other
route. Health checks probe evaluation-mode deployments through the
Decisions API instead of failing with an unsupported mode, and
cloudflare/clef and cloudflare/clef-flash get cost-map rows so the short
names resolve a mode and a price. The registry no longer claims typed
decisions for a provider with no backend.

* fix(decisions): let health_check_params override the evaluation probe

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* test(integration): audit the decisions endpoint across providers, limits, health and chaos

Adds the /v1/decisions audit cells: one wire contract per provider (path, key, body and cost-map billing), the gateway-only fields and tags, the sad paths (invalid bodies, unknown model, key checks, api_base in the body, upstream 401/429/500, a 200 without answers, an unreachable upstream), the two evaluation-mode health probes, and three chaos cells (a mixed-failure burst over both routes, a worker SIGKILL mid-burst, an upstream outage and restart on the same port).

The PR's cost case read the upstream observations through the gateway, which answers 404 for that path; it now reads them from the upstream URL. The owned proxy harness takes extra CLI arguments, and its graceful stop waits as long as a worker boot may take, since a worker still starting honors SIGTERM only once it is up and the 30 second wait forced a cleanup under load.

* fix(decisions): send env API keys to a configured api_base

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* feat(decisions): add zero-cost evaluation cost-map entry for Strands Decider

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* fix(decisions): register the routes through the lazy feature registry

The Decisions router was included at import, ahead of the config and DB
pass-through endpoints, so a pass-through configured at /v1/decisions
was skipped and answered 400 as an unknown Decisions provider. The
routes now register through LAZY_FEATURES, which splices them in after
every eager route, so a pass-through at /v1/decisions keeps its route
while /decisions still serves natively. The lazy OpenAPI snapshot carries
the two paths so the schema shows them before the first call.

The audit cells add the env-key egress to a configured api_base, the
client api_base opt-in shared with chat, the pass-through precedence on
an owned proxy, and the Strands evaluation health check resolved from
the cost map. The integration config exports the Perplexity env key the
first cell needs.

* fix(decisions): keep the Cloudflare api_base message in its transformation and read the audit upstream once per cell

* fix(proxy): let a config pass-through beat a lazily registered route in eager mode

With LITELLM_DISABLE_LAZY_ROUTES set the decisions routes are registered at
startup, so SafeRouteAdder treated a config pass-through at exactly
/v1/decisions as already registered and dropped it. In lazy mode a pass-through
created through the API after the first native call was skipped the same way.
Routes a lazy feature owns no longer count as registered, and a route added at
one of their paths is placed ahead of them, the precedence lazy mode gives a
config pass-through when the feature has not loaded yet.

---------

Co-authored-by: mateo <mateo@berri.ai>
Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
Co-authored-by: mateo-berri <277851410+mateo-berri@users.noreply.github.com>
2026-10-03 17:38:38 +00:00

289 lines
9.2 KiB
Python

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
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(launch: _Launch) -> bool:
return launch.process.poll() is not None and "address already in use" in launch.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)
assert launch.process.poll() is None or (attempts > 1 and _lost_port_race(launch)), (
"Owned proxy exited before readiness"
)
except BaseException:
_stop(launch.process)
raise
if launch.process.poll() is None:
return launch
_stop(launch.process)
return _launch_until_bound(command, root, environment, output, attempts - 1)
@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 = Path(os.environ.get("INTEGRATION_PROXY_ROOT") or Path(__file__).resolve().parents[3])
environment: Final = {
**{
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,
}
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)
_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 = Path(os.environ.get("INTEGRATION_PROXY_ROOT") or Path(__file__).resolve().parents[3])
slot: Final = UpstreamSlot(directory, _free_port(), root)
slot.start()
try:
yield slot
finally:
if slot.process is not None:
slot.stop()