Merge remote-tracking branch 'origin/main' into litellm_tencent_ui_dropdown

This commit is contained in:
Jim Aldon D'Souza 2026-10-01 09:27:55 -07:00
commit 48b653bca1
1603 changed files with 143392 additions and 21088 deletions

View file

@ -17,3 +17,6 @@ rustflags = ["-C", "link-arg=-undefined", "-C", "link-arg=dynamic_lookup"]
[target.aarch64-apple-darwin]
rustflags = ["-C", "link-arg=-undefined", "-C", "link-arg=dynamic_lookup"]
[env]
SQLX_OFFLINE = "true"

View file

@ -3419,7 +3419,7 @@ workflows:
name: integration-<< matrix.suite >>
matrix:
parameters:
suite: [management, accounting, database, providers, mcp, sdk, cost, browser]
suite: [management, accounting, database, providers, mcp, sdk, cost, security, browser]
- integration_contracts:
name: integration-extensions
suite: extensions

View file

@ -46,6 +46,7 @@ legacy_paths() {
echo tests/unit/google_genai
echo tests/unit/router_strategy
echo tests/unit/router_utils
echo tests/unit/proxy/common_utils/test_cache_aware_routing.py
echo tests/unit/enterprise/enterprise_callbacks/send_emails
echo tests/unit/enterprise/proxy/test_afile_retrieve_returns_unified_id.py
echo tests/unit/enterprise/proxy/test_batch_retrieve_input_file_id.py
@ -88,6 +89,7 @@ legacy_paths() {
proxy-db-auth-checks)
echo tests/unit/proxy/auth/test_auth_checks.py
echo tests/unit/proxy/auth/test_user_api_key_auth.py
echo tests/unit/proxy/test_credential_slot_registry.py
echo tests/unit/proxy/test_deprecated_key_grace_period.py ;;
proxy-db-budgets)
echo tests/unit/proxy/auth/test_default_end_user_budget_simple.py
@ -105,6 +107,7 @@ legacy_paths() {
echo tests/unit/proxy/test_update_spend.py
echo tests/unit/skills/test_skills_db.py ;;
proxy-db-endpoints-and-responses)
echo tests/unit/proxy/engine
echo tests/unit/proxy/auth/test_models_fallback_endpoint.py
echo tests/unit/proxy/common_utils/test_check_batch_cost.py
echo tests/unit/proxy/common_utils/test_check_responses_cost.py
@ -146,7 +149,10 @@ legacy_paths() {
echo tests/unit/proxy/test_proxy_server.py ;;
proxy-db-proxy-utils) echo tests/unit/proxy/test_proxy_utils.py ;;
proxy-extras) echo tests/unit/litellm_proxy_extras ;;
proxy-infra) echo tests/unit/gateway ;;
proxy-infra)
echo tests/unit/gateway
echo tests/unit/proxy/management_endpoints/test_roi_calculator_endpoints.py
echo tests/unit/proxy/roi_calculator ;;
responses-caching-types)
find tests/unit/responses -name 'test_*.py' -not -path 'tests/unit/responses/mcp/*'
echo tests/unit/types ;;

View file

@ -42,7 +42,7 @@ case "$subject" in
;;
esac
ALLOWED_TYPES="feat|fix|docs|style|refactor|perf|test|build|ci|chore|revert"
ALLOWED_TYPES="feat|fix|docs|style|refactor|perf|test|build|ci|chore|revert|security"
# Description must not start with an uppercase letter — kept in sync with the
# subjectPattern in .github/workflows/conventional-commits.yml so the local
# hook is the strictly tighter of the two gates. (Without this guard, a commit
@ -61,7 +61,7 @@ cat >&2 <<EOF
Expected: <type>(<scope>)!: <description>
(description must start with a lowercase letter)
Allowed types: feat, fix, docs, style, refactor, perf, test, build, ci, chore, revert
Allowed types: feat, fix, docs, style, refactor, perf, test, build, ci, chore, revert, security
Examples:
feat(router): add weighted round-robin strategy
fix(bedrock): decouple STS region from aws_region_name

8
.github/CODEOWNERS vendored
View file

@ -1,10 +1,2 @@
/ui/ @yuneng-berri @ryan-crabbe-berri
/litellm/proxy/_experimental/out/ @yuneng-berri @ryan-crabbe-berri
/ui/Dockerfile
/ui/nginx.conf
/ui/litellm-dashboard/src/lib/http/schema.d.ts
/ui/litellm-dashboard/tsconfig.tsbuildinfo
/model_prices_and_context_window.json @mateo-berri @ryan-crabbe-berri @kerry-berri
/litellm/model_prices_and_context_window_backup.json @mateo-berri @ryan-crabbe-berri @kerry-berri
/litellm-proxy-extras/litellm_proxy_extras/migrations/ @yuneng-berri @ryan-crabbe-berri
/.github/CODEOWNERS @yuneng-berri

Binary file not shown.

After

Width:  |  Height:  |  Size: 80 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 58 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 63 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 70 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 47 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 76 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 75 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 39 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 70 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 93 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 73 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 39 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 81 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 72 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 76 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 50 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 59 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 57 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 56 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 65 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 55 KiB

View file

@ -111,3 +111,10 @@ dockerfiles:
An example image under cookbook/ that is documentation rather than a shipped artifact
paths:
- cookbook/litellm-ollama-docker-image/Dockerfile
- reason: >-
The Rust gateway image compiles the whole workspace in release mode, which is too slow for
a per-pull-request job while the gateway binary is still being assembled; the Rust lint,
clippy, and compile jobs already cover the code it packages. Revisit when the gateway is
published
paths:
- litellm-rust/crates/gateway/Dockerfile

View file

@ -214,7 +214,7 @@ def main(
native_module: Final = load_native_module(native_path)
native_module_loads: Final = native_module is not None
panic_test_hook_absent: Final = native_module is not None and not hasattr(native_module, "_panic_for_test")
native_size_limit: Final = 40_000_000
native_size_limit: Final = 45_000_000
native_size_within_limit: Final = native_member.file_size <= native_size_limit
validations: Final = (
(f"Python tag is {EXPECTED_PYTHON_TAG}", python_tag == EXPECTED_PYTHON_TAG),

View file

@ -41,6 +41,7 @@ jobs:
ci
chore
revert
security
requireScope: false
subjectPattern: ^(?![A-Z]).+$
subjectPatternError: |

View file

@ -15,6 +15,8 @@ jobs:
runs-on: ubuntu-latest
permissions:
contents: write
outputs:
version: ${{ steps.version.outputs.version }}
steps:
- name: Require main
env:
@ -64,3 +66,14 @@ jobs:
sha: context.sha,
});
core.info(`Created branch ${branchName} at ${context.sha}`);
linear-release:
name: Move the Linear release to rc
needs: create-rc-branch
permissions:
contents: read
uses: ./.github/workflows/linear-release.yml
with:
rc_version: ${{ needs.create-rc-branch.outputs.version }}
secrets:
LINEAR_API_KEY: ${{ secrets.LINEAR_API_KEY }}

60
.github/workflows/lens-worker.yml vendored Normal file
View file

@ -0,0 +1,60 @@
name: Lens Worker Image
on:
pull_request:
branches: [main, litellm_oss_branch, "litellm_**"]
paths:
- deploy/lens/**
- litellm/proxy/engine/**
- .github/workflows/lens-worker.yml
push:
branches: [main, litellm_agent_engine]
paths:
- deploy/lens/**
- litellm/proxy/engine/**
- .github/workflows/lens-worker.yml
workflow_dispatch:
permissions:
contents: read
concurrency:
group: ${{ github.workflow }}-${{ github.event.pull_request.number || github.ref }}
cancel-in-progress: true
jobs:
lens-worker-image:
permissions:
contents: read
packages: write
runs-on: ubuntu-latest
timeout-minutes: 10
steps:
- uses: actions/checkout@08eba0b27e820071cde6df949e0beb9ba4906955 # v4.3.0
with:
persist-credentials: false
- name: Build Lens worker
run: docker build -f deploy/lens/Dockerfile -t lens-worker:${{ github.sha }} .
- name: Verify standalone imports with a read-only filesystem
run: |
docker run --rm --network none --read-only --cap-drop ALL \
--security-opt no-new-privileges --entrypoint python \
lens-worker:${{ github.sha }} -c '
import os
import engine.worker
from engine.trace_store import trace_store
assert os.getuid() == 65532
with trace_store() as store:
assert store.count() == 0
'
- name: Publish versioned Lens worker
if: github.event_name != 'pull_request' && github.repository == 'BerriAI/litellm'
env:
REGISTRY_TOKEN: ${{ secrets.GITHUB_TOKEN }}
REGISTRY_USER: ${{ github.actor }}
IMAGE: ghcr.io/berriai/litellm-lens-worker:sha-${{ github.sha }}
run: |
printf '%s' "$REGISTRY_TOKEN" | docker login ghcr.io -u "$REGISTRY_USER" --password-stdin
docker tag lens-worker:${{ github.sha }} "$IMAGE"
docker push "$IMAGE"
printf 'Lens worker image: `%s`\n' "$IMAGE" >> "$GITHUB_STEP_SUMMARY"

131
.github/workflows/linear-release.yml vendored Normal file
View file

@ -0,0 +1,131 @@
name: Linear Release
on:
push:
branches:
- main
- "rc/**"
release:
types: [published]
workflow_call:
inputs:
rc_version:
description: "X.Y.0 release whose rc branch was just cut"
required: true
type: string
secrets:
LINEAR_API_KEY:
required: true
permissions: {}
jobs:
linear-release:
name: Linear Release
if: github.repository == 'BerriAI/litellm'
runs-on: ubuntu-latest
permissions:
contents: read
steps:
- uses: actions/checkout@08eba0b27e820071cde6df949e0beb9ba4906955 # v4.3.0
with:
fetch-depth: 0
persist-credentials: false
- name: Plan
id: plan
env:
EVENT: ${{ github.event_name }}
REF_NAME: ${{ github.ref_name }}
BEFORE: ${{ github.event.before }}
CREATED: ${{ github.event.created }}
RC_VERSION: ${{ inputs.rc_version }}
RELEASE_TAG: ${{ github.event.release.tag_name }}
PRERELEASE: ${{ github.event.release.prerelease }}
run: |
set -euo pipefail
sync_base="${BEFORE}"
if [ "${CREATED}" = "true" ]; then
sync_base=""
fi
if [ -n "${RC_VERSION}" ]; then
echo "version=${RC_VERSION}" >> "$GITHUB_OUTPUT"
echo "stage=rc" >> "$GITHUB_OUTPUT"
elif [ "${EVENT}" = "release" ]; then
if [ "${PRERELEASE}" = "true" ] || ! echo "${RELEASE_TAG}" | grep -qE '^v[0-9]+\.[0-9]+\.0$'; then
echo "::notice::${RELEASE_TAG} is not an X.Y.0 stable release; nothing to complete"
exit 0
fi
echo "version=${RELEASE_TAG#v}" >> "$GITHUB_OUTPUT"
echo "complete=true" >> "$GITHUB_OUTPUT"
elif [ "${REF_NAME}" = "main" ]; then
version="$(python3 .github/scripts/read_rc_version.py | cut -d= -f2)"
status=0
git ls-remote --exit-code --heads origin "rc/${version}" > /dev/null || status=$?
case "${status}" in
0)
IFS=. read -r major minor _ <<< "${version}"
version="${major}.$((minor + 1)).0"
;;
2) ;;
*)
echo "::error::could not check whether rc/${version} exists (git ls-remote exit ${status})"
exit 1
;;
esac
echo "version=${version}" >> "$GITHUB_OUTPUT"
echo "sync_base=${sync_base}" >> "$GITHUB_OUTPUT"
echo "main=true" >> "$GITHUB_OUTPUT"
else
echo "version=${REF_NAME#rc/}" >> "$GITHUB_OUTPUT"
echo "sync_base=${sync_base}" >> "$GITHUB_OUTPUT"
echo "stage=rc" >> "$GITHUB_OUTPUT"
fi
- name: Sync commits into the release
if: steps.plan.outputs.sync_base != ''
uses: linear/linear-release-action@d4af10092984f9bc6d5efa075b242bdf01333463 # v0.18.0
with:
access_key: ${{ secrets.LINEAR_API_KEY }}
command: sync
name: LiteLLM ${{ steps.plan.outputs.version }}
version: ${{ steps.plan.outputs.version }}
base_ref: ${{ steps.plan.outputs.sync_base }}
cli_version: v0.18.0
- name: Keep the main stage unless the rc branch was cut during this run
id: main_stage
if: steps.plan.outputs.main == 'true'
env:
VERSION: ${{ steps.plan.outputs.version }}
run: |
set -euo pipefail
status=0
git ls-remote --exit-code --heads origin "rc/${VERSION}" > /dev/null || status=$?
case "${status}" in
0) echo "::notice::rc/${VERSION} was cut during this run; leaving the release in its rc stage" ;;
2) echo "stage=main" >> "$GITHUB_OUTPUT" ;;
*)
echo "::error::could not check whether rc/${VERSION} exists (git ls-remote exit ${status})"
exit 1
;;
esac
- name: Move the release to its stage
if: steps.plan.outputs.stage != '' || steps.main_stage.outputs.stage != ''
uses: linear/linear-release-action@d4af10092984f9bc6d5efa075b242bdf01333463 # v0.18.0
with:
access_key: ${{ secrets.LINEAR_API_KEY }}
command: update
stage: ${{ steps.plan.outputs.stage || steps.main_stage.outputs.stage }}
version: ${{ steps.plan.outputs.version }}
cli_version: v0.18.0
- name: Complete the release
if: steps.plan.outputs.complete == 'true'
uses: linear/linear-release-action@d4af10092984f9bc6d5efa075b242bdf01333463 # v0.18.0
with:
access_key: ${{ secrets.LINEAR_API_KEY }}
command: complete
version: ${{ steps.plan.outputs.version }}
cli_version: v0.18.0

View file

@ -80,6 +80,11 @@ jobs:
- name: test_e2e_changed_gate
run: uv run --no-sync pytest -q --noconftest -p no:cacheprovider -c /dev/null tests/code_coverage_tests/test_e2e_changed_gate.py tests/code_coverage_tests/test_e2e_idp_stack.py
- name: test_e2e_metadata
env:
PYTHONPATH: tests/e2e
run: uv run --no-sync pytest -q --noconftest -p no:cacheprovider -c /dev/null tests/code_coverage_tests/test_e2e_metadata.py tests/code_coverage_tests/test_e2e_junit_report.py
- name: Check merge smoke harness
run: uv run --no-sync pytest -q --noconftest -p no:cacheprovider -c /dev/null tests/code_coverage_tests/test_merge_smoke.py

View file

@ -24,6 +24,7 @@ jobs:
timeout-minutes: ${{ matrix.job-timeout-minutes }}
permissions:
contents: read
id-token: write
services:
postgres:
@ -134,9 +135,19 @@ jobs:
env:
TEST_PATH: ${{ matrix.test-path }}
WORKERS: ${{ matrix.workers }}
PYTEST_ADDOPTS: ${{ matrix.shard == 'proxy-behavior' && '--cov=litellm/proxy/engine --cov-report=xml:coverage-lens-postgres.xml' || '' }}
run: |
if [ "${WORKERS}" = "0" ]; then
uv run --no-sync pytest ${TEST_PATH:?} -vv --tb=short --durations=10
else
uv run --no-sync pytest ${TEST_PATH:?} -vv --tb=short --durations=10 -n "${WORKERS}"
fi
- name: Upload Lens database coverage
if: steps.changes.outputs.decision != 'skip' && matrix.shard == 'proxy-behavior' && !cancelled()
uses: codecov/codecov-action@75cd11691c0faa626561e295848008c8a7dddffe # v5.5.4
with:
use_oidc: true
files: coverage-lens-postgres.xml
flags: lens-postgres
fail_ci_if_error: true

View file

@ -79,7 +79,9 @@ jobs:
- shard: integrations
artifact-name: integrations
test-path: ""
test-path: >-
tests/test_litellm/integrations
tests/test_litellm/tracing
unit-flag: integrations
workers: 2
reruns: 3

View file

@ -4,7 +4,7 @@
.PHONY: help test test-unit test-unit-llms test-unit-proxy-guardrails test-unit-proxy-core test-unit-proxy-misc \
test-unit-integrations test-unit-core-utils test-unit-other test-unit-root \
test-proxy-unit-a test-proxy-unit-b test-integration test-unit-helm \
test-rust-extension \
test-rust-extension rust-sqlx-prepare \
info lint lint-inner lint-dev lint-checks format \
lint-basedpyright lint-e2e-basedpyright lint-basedpyright-budget-update lint-type-discipline lint-type-discipline-budget-update \
lint-ruff-budget lint-ruff-budget-update lint-budget-update lint-gate \
@ -56,6 +56,7 @@ help:
@echo " make test-integration - Run integration tests"
@echo " make test-unit-helm - Run helm unit tests"
@echo " make test-rust-extension - Build the Rust extension and run its public Python tests"
@echo " make rust-sqlx-prepare - Refresh litellm-rust/crates/db/.sqlx against a migrated Postgres container"
@echo ""
@echo "Heavy targets (check, lint) queue for LITELLM_GATE_SLOTS machine-wide"
@echo "slots (default 2; 0 disables) so parallel sessions don't thrash one machine."
@ -306,6 +307,9 @@ test-rust-extension:
LITELLM_RUST=1 LITELLM_LOCAL_MODEL_COST_MAP=True \
"$$temporary/venv/bin/python" -I -m pytest --import-mode=importlib -m requires_rust_extension tests/test_litellm_rust
rust-sqlx-prepare:
cd litellm-rust && cargo run -p litellm-db-testing --bin sqlx-prepare
test: install-test-deps
$(UV_RUN) pytest tests/

View file

@ -81,6 +81,8 @@ BACKEND_PATH_PREFIXES: tuple[str, ...] = (
# Spend / analytics
"/spend/",
"/analytics/",
"/engine/",
"/v1/traces",
"/global/",
"/user_agent",
"/usage/",
@ -144,6 +146,7 @@ BACKEND_EXACT_PATHS: frozenset[str] = frozenset(
{
"/",
"/routes",
"/engine",
"/openapi.json",
"/docs",
"/docs/oauth2-redirect",

View file

@ -99,7 +99,7 @@
"limit": 0
},
"reportUnknownArgumentType": {
"limit": 44358
"limit": 44802
},
"reportUnknownLambdaType": {
"limit": 109

View file

@ -1,37 +0,0 @@
# Publish MCP servers in the AI Hub
Set `litellm_settings.public_mcp_servers` to the concrete IDs of the servers you want listed in the public AI Hub. Pin `server_id` in each configuration entry so the publication list stays stable across deployments
```yaml
mcp_servers:
documentation:
server_id: documentation-mcp
url: https://mcp.example.com/mcp
transport: http
available_on_public_internet: true
litellm_settings:
public_mcp_hub_strict_whitelist: true
public_mcp_servers:
- documentation-mcp
```
Use `documentation-mcp`, the `server_id`, in the publication list. The configuration key `documentation`, display names, and aliases are not publication IDs. Database-created servers use the ID returned by `/v1/mcp/server`
The dashboard's **AI Hub > MCP Hub > Manage MCP Hub Visibility** dialog edits this same list. Its YAML example includes the selected server IDs. With database-backed configuration (`store_model_in_db: true`), a value declared in YAML is owned by that file: edit the file and reload, or remove that key from YAML to let the dashboard manage it in the database. File-backed deployments can save the list directly to their configuration file
To remove all explicit entries, save an empty selection in the dialog or configure:
```yaml
litellm_settings:
public_mcp_hub_strict_whitelist: true
public_mcp_servers: []
```
## Hub listing and network access
The **Hub listing** column in AI Hub identifies servers that appear in `/public/mcp_hub`. The dashboard derives this status from the current registry and publication settings. Setting `mcp_info.is_public` on a server does not publish it; that response field is derived metadata. `mcp_info.is_public_explicit` identifies registered servers included in the explicit publication list
Gateway cards and server details show **All Networks** when `available_on_public_internet` is enabled or the server is explicitly published in `public_mcp_servers`. They show **Internal Only** when both are false. The per-server flag defaults to `true`; explicit publication overrides a disabled flag for compatibility. Older proxies that omit the metadata needed to determine access show **Unknown**. These labels describe allowed client IPs; authentication and tool permissions still apply
The default `public_mcp_hub_strict_whitelist: true` lists only registered servers in `public_mcp_servers`. Legacy mode (`false`) additionally lists registered servers with `available_on_public_internet: true`. In legacy mode, clearing the explicit publication list leaves these automatically listed servers visible. Enable strict mode when the publication list should fully determine hub visibility

View file

@ -24,7 +24,7 @@ model_list:
- model_name: sagemaker-completion-model
litellm_params:
model: sagemaker/berri-benchmarking-Llama-2-70b-chat-hf-4
input_cost_per_second: 0.000420
cost_per_second: 0.000420
- model_name: text-embedding-ada-002
litellm_params:
model: azure/azure-embedding-model

View file

@ -12,7 +12,7 @@ def encode_image(image_path):
# Path to your image
image_path = "litellm/proxy/logo.jpg"
image_path = "litellm/proxy/logo.png"
# Getting the Base64 string
base64_image = encode_image(image_path)
@ -27,7 +27,7 @@ response = client.responses.create(
{"type": "input_text", "text": "what color is the image"},
{
"type": "input_image",
"image_url": f"data:image/jpeg;base64,{base64_image}",
"image_url": f"data:image/png;base64,{base64_image}",
},
],
}

7
deploy/lens/Dockerfile Normal file
View file

@ -0,0 +1,7 @@
FROM python:3.12-slim
WORKDIR /app
RUN pip install --no-cache-dir httpx==0.28.1 pydantic==2.11.7
COPY litellm/proxy/engine/__init__.py litellm/proxy/engine/models.py litellm/proxy/engine/trace_store.py litellm/proxy/engine/analysis.py litellm/proxy/engine/worker.py /app/engine/
VOLUME /tmp
USER 65532:65532
CMD ["python", "-m", "engine.worker"]

View file

@ -0,0 +1,8 @@
**
!litellm/
!litellm/proxy/
!litellm/proxy/engine/
!litellm/proxy/engine/__init__.py
!litellm/proxy/engine/models.py
!litellm/proxy/engine/analysis.py
!litellm/proxy/engine/worker.py

105
deploy/lens/README.md Normal file
View file

@ -0,0 +1,105 @@
# Lens worker
Lens reviews recorded activity and saves evidence-linked findings in the LiteLLM dashboard under Observability, Lens (`/ui/lens/`)
## Start a worker
Upgrade your existing LiteLLM proxy to a release that includes Lens with PostgreSQL, agent tracing (`general_settings.tracing: {store: clickhouse}`), and ClickHouse configured through `CLICKHOUSE_URL` and a separate SELECT-only `CLICKHOUSE_READER_URL`. Enable the ClickHouse callback and request/response logging to analyze LLM requests. Lens can only inspect content you actually retain
In Lens, click **Connect worker**, then **Generate setup command**. The LiteLLM address is filled in for you; change it only if the server running Docker needs a different network address. Copy the command and run it on your server. The dialog changes to **Worker connected** when the container checks in
The command already contains the compatible worker image and one worker token. No separate API key, source checkout, environment file, or second LiteLLM deployment is needed. Keep the command private because it includes the token. The LiteLLM release provides the dashboard and APIs; the container only runs background analysis
The dashboard and Compose file pin a verified worker image by digest. The image uses Linux amd64, and the generated command selects that platform. Worker image releases are independent of proxy releases: update the pinned image when changing their API contract. CI also publishes immutable commit tags for reproducible builds
For deployments managed with Compose, download `compose.yaml` and provide `LITELLM_URL` and `LENS_WORKER_TOKEN` in an environment file. Its default image is already selected:
```bash
docker compose --env-file /path/to/lens.env -f compose.yaml up -d
```
Developers can build locally with `LENS_WORKER_IMAGE=litellm-lens-worker:local docker compose -f deploy/lens/compose.yaml -f deploy/lens/compose.build.yaml up -d --build`
The worker needs outbound HTTPS access to LiteLLM. It needs no inbound ports, provider keys, direct database access, or GPU. The proxy calls your selected model through its configured router; trace content reaches that model provider. Use a model with JSON output support and known token prices. One worker handles one scan at a time and can serve multiple lenses. For more throughput, start another worker with a separate credential
V1 setup, manual runs, feedback, and worker credentials are restricted to proxy administrators. Proxy-admin viewers can inspect results. Regular user and team keys cannot access the Lens API. Worker credentials can serve the administrator’s lenses. Revoke it in the connection dialog when retiring a worker. Redeploy the worker alongside proxy upgrades so their API versions match
## Configure a lens
Choose agent runs, individual LLM requests, or both. The matching-activity preview updates as you choose an application (the recorded OpenTelemetry service.name) or, for request activity, a LiteLLM model group and add metadata conditions. It shows run names, timestamps, and trace IDs; open a run to inspect its original steps before starting analysis. Suggestions come from up to 100 recent executions and may not include every recorded attribute. You can enter other exact keys and values. Leave service and filters blank for all activity your account can access. Filters are exact key/value matches, combined with AND. Trace filters match span or resource attributes on the same span. Request filters match logged metadata, including caller metadata stored under `requester_metadata`; `tag=value` matches request tags. `swarm=research` works only if your instrumentation records that attribute
Describe how the agent should behave and optionally add specific checks. Select the lookback window, team and metadata, then choose the percentage to review and an optional maximum. **100% with no maximum selects every matching run**. The preview pages through all matching activity and lets you select particular runs. Percentage sampling uses a stable hash order, rounds up, and applies the optional maximum after the percentage
Choose your analysis model, parallelism and monthly budget. Parallelism controls simultaneous model calls, not the number of runs selected. New lenses run once by default. Turn on monitoring to repeat the same setup at a custom interval. **Run now** uses the same saved settings immediately, including the same lookback window and sampling. Every scan recalculates the window, so overlapping windows can review the same activity again. Duplicate a lens when you want a separate investigation without changing an existing monitor
Pausing stops future scheduled scans; cancel the active scan separately if needed. The worker polls every 10 seconds; creating a lens or clicking Run now queues a scan, and due schedules are queued when the worker polls. Scans for the same lens never overlap, and its next interval starts after completion. Closing the browser does not stop the worker. Configuration edits apply to the next scan. A running scan retains its settings and selected execution IDs across retries
## Read the results
Needs attention shows issues, highest priority first. Patterns contains useful trends and successful behavior that may not need a fix. Each finding starts with a short explanation and a next step when useful. Expand the limitations for uncertainty and counterexamples. Evidence is grouped by run and collapsed until you need it; each quote opens the original step
Use the batch selector or Scans tab to reopen previous results. Each batch keeps its own findings, settings, selected runs, coverage and cost. Older batches created before snapshot support remain available through accumulated findings. The Runs tab lists the selected batch's sample and can filter per-run observations, including runs without an observed issue and runs with insufficient evidence. These observations precede the final evidence investigation. Linked-run counts on findings include cited counterexamples, so they are not failure counts
Choose **This is expected** and explain why to teach later scans about acceptable behavior. Feedback is kept with the lens and included in subsequent reviews. It does not alter historical evidence or exempt different problems
## What a scan does
The proxy selects executions received or updated within the configured lookback window, with a two-minute settling period. Older rows without receipt timestamps use execution end time. Overlapping scans do not increment a finding's occurrence count for the same execution ID
A trace is spans sharing a trace ID within one team, not an automatically reconstructed conversation session. Requests are individual LLM calls. When both sources are enabled, requests correlated to a recorded span by response ID are excluded to reduce double counting
The worker reviews the selected executions in parallel. It pages through their recorded spans and gives the first reviewer a catalog, task and outcome excerpts. The reviewer can read more original content to resolve uncertainties. Large catalogs and groups of observations are processed in bounded context windows, with every page available. Grouping retains supporting run IDs in code, so a pattern occurring thousands of times does not require a model to repeat thousands of IDs. Candidate investigators can page through supporting observations, other runs and original evidence
There is no fixed total run, span, candidate or investigation-turn cutoff. Repeated or empty evidence requests stop a stalled investigation. Context windows, the configured budget, available model capacity and recorded evidence still bound practical work. The dashboard reports completed work and gaps. The investigator has no shell, browsing, code-editing or production-action tools
Each model response must match a bounded JSON schema. A malformed response gets one repair attempt through the same budget controls; repeated invalid output fails the scan. Both the worker and proxy validate quoted evidence. Findings retain exact quotes and open the source trace or request. Resolve a finding after a fix, or dismiss it with a reason. A resolved finding reopens when new execution IDs support the same pattern; dismissed findings remain dismissed
Coverage distinguishes eligible, sampled, reviewed, partial, and unassessable executions. Findings describe observations in the sample, not population-wide success rates or proven causes. A root span does not prove that a trace contains every expected span. Long, missing, redacted, or expired content limits the conclusions
## Operations and limits
PostgreSQL stores configurations, findings and all scan history, returned in pages of 50 jobs. Workers claim jobs with optimistic concurrency and a five-minute lease, renewed every 30 seconds. A disconnected job can be reclaimed up to three times. Cancellation stops subsequent work; a model call already in flight may finish and incur cost
Before every model call, Lens reserves a conservative amount against the monthly lens budget. Successful calls reconcile to reported cost where pricing is available. Interrupted calls retain their reservation because the provider may have charged. A scan stops when the next reservation would exceed the limit, so it can stop with some budget remaining. Lens budgets are separate from virtual-key budgets; analysis calls use the proxy router directly
V1 requires ClickHouse for both sources. It does not reconstruct sessions from unrelated trace IDs, guarantee exhaustive reviews, cache all per-execution observations across scans, or automatically fix agent code. Trace contents can change as late spans arrive, even though a job's selected IDs are fixed. Findings should be reviewed by a person before acting on them
## API access
The UI and API use the same scan lifecycle. Authenticate with a proxy administrator credential for writes, or a proxy-admin viewer credential for reads. Worker credentials are only for worker operations
```bash
curl "$LITELLM_URL/engine" -H "Authorization: Bearer $LITELLM_API_KEY" \
-H 'Content-Type: application/json' -d '{
"name": "Research quality", "model": "your-model-alias",
"context": "Answer the requested question using cited, retrieved evidence.",
"source": "traces", "lookback_hours": 24,
"sample_percent": 100, "sample_size": null, "concurrency": 8,
"enabled": true, "interval_minutes": 1440, "monthly_budget": 50
}'
curl "$LITELLM_URL/engine/$LENS_ID/runs" -X POST \
-H "Authorization: Bearer $LITELLM_API_KEY" -H 'Content-Type: application/json' -d '{}'
curl "$LITELLM_URL/engine/$LENS_ID/runs?offset=0" -H "Authorization: Bearer $LITELLM_API_KEY"
curl "$LITELLM_URL/engine/$LENS_ID/runs/$BATCH_ID" -H "Authorization: Bearer $LITELLM_API_KEY"
```
Creation queues the first batch. Posting to `/engine/{id}/runs` queues another, or returns the existing active batch. The run response contains its ID under `jobs[0].id`. Poll the batch URL for status, findings and assessments. List responses omit large result payloads; request a batch to retrieve them. Supply an optional complete `settings` object on the runs POST for a one-off override; the saved lens stays unchanged. Selection accepts `team_id`, exact `filters`, and opaque `execution_ids` returned by `/engine/preview/sample`. Preview accepts `offset` and `as_of` to keep the time window fixed while paging. Feedback uses `PATCH /engine/{id}/findings/{finding_id}` with `status` and `reason`
## Quality evaluation
Run the checked-in cases against a configured real model. Expected labels are used only for scoring, never passed to the model. Dev and held-out cases include missing outcomes, failed tools, recovery, handoffs, unsupported claims, repeated work, long evidence and prompt injection. The background option adds clean arithmetic traces to test rare-issue discovery at scale; those repeated synthetic cases do not establish accuracy on every production workload
```bash
python -m tests.proxy_behavior.lens.evaluate --api-base "$LITELLM_URL" \
--model your-model-alias --split all --background 1000 --concurrency 16 \
--output /tmp/lens-quality.json
```
Set `LITELLM_API_KEY` privately. This makes paid model calls. Inspect missed and unexpected per-run labels, final findings and coverage; do not equate a passing dataset with guaranteed detection on arbitrary traces
The worker uses temporary disk space for trace content while reviewing it, and removes those files after each review. Its Docker image supplies a writable temporary volume while keeping the application filesystem read-only
To check that accepted behavior stays accepted without hiding new problems, run the evaluator with `--dataset tests/proxy_behavior/lens/feedback_cases.json`. Reports include elapsed time, model call count, reported cost when the proxy provides it, missed checks, unexpected checks, and inconclusive candidates

View file

@ -0,0 +1,6 @@
services:
lens-worker:
build:
context: ../..
dockerfile: deploy/lens/Dockerfile
image: litellm-lens-worker:local

10
deploy/lens/compose.yaml Normal file
View file

@ -0,0 +1,10 @@
services:
lens-worker:
image: ${LENS_WORKER_IMAGE:-ghcr.io/berriai/litellm-lens-worker@sha256:40fdb82113dd4474cb6e833cf28552487d87c8baf61693a1c3fc2863b7968c6a}
environment:
LITELLM_URL: ${LITELLM_URL:?Set the URL reachable from this container}
LENS_WORKER_TOKEN: ${LENS_WORKER_TOKEN:?Create a worker credential in the Lens UI}
restart: unless-stopped
read_only: true
cap_drop: [ALL]
security_opt: [no-new-privileges:true]

Binary file not shown.

After

Width:  |  Height:  |  Size: 95 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 6.9 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 89 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 80 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 70 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 132 KiB

View file

@ -103,7 +103,7 @@ ENV LITELLM_NON_ROOT=true
RUN mkdir -p /var/lib/litellm/ui /var/lib/litellm/assets && \
cp -r /app/litellm/proxy/_experimental/out/. /var/lib/litellm/ui/ && \
cp /app/litellm/proxy/logo.jpg /var/lib/litellm/assets/logo.jpg && \
cp /app/litellm/proxy/logo.png /var/lib/litellm/assets/logo.png && \
touch /var/lib/litellm/ui/.litellm_ui_ready
RUN --mount=type=cache,target=/app/.cache/uv,id=litellm-uv-cache \

View file

@ -0,0 +1,62 @@
name: litellm-tracing
services:
litellm:
build:
context: ..
target: runtime
command: ["--config", "/app/tracing-config.yaml", "--port", "4000"]
environment:
LITELLM_MASTER_KEY: local-tracing-master-key
LITELLM_SALT_KEY: sk-local-tracing-salt-key
DATABASE_URL: postgresql://litellm:litellm@db:5432/litellm
STORE_MODEL_IN_DB: "True"
CLICKHOUSE_URL: http://default:local-tracing@clickhouse:8123
CLICKHOUSE_READER_URL: http://default:local-tracing@clickhouse:8123
CLICKHOUSE_DATABASE: litellm
OPENAI_API_KEY: ${OPENAI_API_KEY:-}
volumes:
- ./tracing-config.yaml:/app/tracing-config.yaml:ro
ports:
- "127.0.0.1:4002:4000"
depends_on:
db:
condition: service_healthy
clickhouse:
condition: service_healthy
db:
image: postgres:16
environment:
POSTGRES_DB: litellm
POSTGRES_USER: litellm
POSTGRES_PASSWORD: litellm
volumes:
- postgres_data:/var/lib/postgresql/data
ports:
- "127.0.0.1:15432:5432"
healthcheck:
test: ["CMD-SHELL", "pg_isready -U litellm -d litellm"]
interval: 5s
timeout: 5s
retries: 10
clickhouse:
image: clickhouse/clickhouse-server:26.9.6.6
environment:
CLICKHOUSE_USER: default
CLICKHOUSE_PASSWORD: local-tracing
CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT: "1"
volumes:
- clickhouse_data:/var/lib/clickhouse
ports:
- "127.0.0.1:18123:8123"
healthcheck:
test: ["CMD", "clickhouse-client", "--user", "default", "--password", "local-tracing", "--query", "SELECT 1"]
interval: 5s
timeout: 5s
retries: 20
volumes:
postgres_data:
clickhouse_data:

View file

@ -0,0 +1,10 @@
model_list:
- model_name: gpt-6.1-sol
litellm_params:
model: openai/gpt-6.1-sol
api_key: os.environ/OPENAI_API_KEY
general_settings:
master_key: os.environ/LITELLM_MASTER_KEY
tracing:
store: clickhouse

View file

@ -22,10 +22,8 @@ from litellm._uuid import uuid
from litellm.proxy._types import *
from litellm.proxy.auth.auth_checks import delete_cached_project_object
from litellm.proxy.auth.user_api_key_auth import user_api_key_auth
from litellm.proxy.management_endpoints.common_utils import (
_is_user_team_admin, # pyright: ignore[reportPrivateUsage] # shared owner of team-admin membership
_set_object_metadata_field,
)
from litellm.proxy.management.teams.access import is_team_admin
from litellm.proxy.management_endpoints.common_utils import _set_object_metadata_field
from litellm.proxy.management_endpoints.team_admin_field_permissions import team_admin_may_manage_projects
from litellm.proxy.management_helpers.utils import (
management_endpoint_wrapper,
@ -117,7 +115,7 @@ async def _check_user_permission_for_project(
return False
team: Final = LiteLLM_TeamTable.model_validate(team_row.model_dump())
return _is_user_team_admin(user_api_key_dict, team) or user_api_key_dict.user_id in (team.admins or [])
return is_team_admin(user_api_key_dict, team) or user_api_key_dict.user_id in (team.admins or [])
async def _validate_team_exists(

View file

@ -1,6 +1,6 @@
[project]
name = "litellm-enterprise"
version = "0.1.71"
version = "0.1.72"
description = "Package for LiteLLM Enterprise features"
readme = "README.md"
requires-python = ">=3.9"
@ -26,7 +26,7 @@ required-version = ">=0.10.9"
module-root = ""
[tool.commitizen]
version = "0.1.71"
version = "0.1.72"
version_files = [
"pyproject.toml:^version",
"../pyproject.toml:litellm-enterprise==",

View file

@ -73,6 +73,7 @@ GATEWAY_PATH_PREFIXES: tuple[str, ...] = (
"/v1/containers",
"/containers",
"/v1/evals",
"/v1/traces",
"/v1/memory",
"/queue/chat/",
# Google data plane (v1beta is the Google AI Studio version)

View file

@ -0,0 +1,97 @@
-- AlterTable
ALTER TABLE "LiteLLM_AgentsTable" ADD COLUMN IF NOT EXISTS "enabled" BOOLEAN NOT NULL DEFAULT true,
ADD COLUMN IF NOT EXISTS "execution_mode" TEXT NOT NULL DEFAULT 'autonomous',
ADD COLUMN IF NOT EXISTS "identity_managed" BOOLEAN NOT NULL DEFAULT false;
-- AlterTable
ALTER TABLE "LiteLLM_SpendLogs" ADD COLUMN IF NOT EXISTS "billing_agent_id" TEXT;
-- CreateTable
CREATE TABLE IF NOT EXISTS "LiteLLM_AgentIdentity" (
"agent_id" TEXT NOT NULL,
"active" BOOLEAN NOT NULL DEFAULT true,
"provider" TEXT NOT NULL,
"issuer" TEXT NOT NULL,
"tenant_id" TEXT NOT NULL,
"client_id" TEXT NOT NULL,
"service_principal_id" TEXT,
"required_roles" TEXT[] DEFAULT ARRAY[]::TEXT[],
"required_scopes" TEXT[] DEFAULT ARRAY['user_impersonation']::TEXT[],
"revision" TEXT NOT NULL,
"last_authenticated_at" TIMESTAMP(3),
CONSTRAINT "LiteLLM_AgentIdentity_pkey" PRIMARY KEY ("agent_id")
);
-- CreateTable
CREATE TABLE IF NOT EXISTS "LiteLLM_RetiredAgentIdentity" (
"binding_id" TEXT NOT NULL,
"agent_id" TEXT,
"provider" TEXT NOT NULL,
"issuer" TEXT NOT NULL,
"tenant_id" TEXT NOT NULL,
"client_id" TEXT NOT NULL,
CONSTRAINT "LiteLLM_RetiredAgentIdentity_pkey" PRIMARY KEY ("binding_id")
);
-- CreateTable
CREATE TABLE IF NOT EXISTS "LiteLLM_RetiredAgent" (
"original_agent_id" TEXT NOT NULL,
"retired_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
CONSTRAINT "LiteLLM_RetiredAgent_pkey" PRIMARY KEY ("original_agent_id")
);
-- CreateTable
CREATE TABLE IF NOT EXISTS "LiteLLM_VerifiedSubject" (
"subject_id" TEXT NOT NULL,
"issuer" TEXT NOT NULL,
"tenant_id" TEXT NOT NULL,
"oid" TEXT NOT NULL,
"kind" TEXT NOT NULL DEFAULT 'human',
"user_id" TEXT,
"verified_via" TEXT NOT NULL DEFAULT 'sso_interactive',
"verified_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
CONSTRAINT "LiteLLM_VerifiedSubject_pkey" PRIMARY KEY ("subject_id")
);
-- CreateIndex
CREATE UNIQUE INDEX IF NOT EXISTS "LiteLLM_AgentIdentity_provider_tenant_id_client_id_key" ON "LiteLLM_AgentIdentity"("provider", "tenant_id", "client_id");
-- CreateIndex
CREATE UNIQUE INDEX IF NOT EXISTS "LiteLLM_AgentIdentity_issuer_service_principal_id_key" ON "LiteLLM_AgentIdentity"("issuer", "service_principal_id");
-- CreateIndex
CREATE UNIQUE INDEX IF NOT EXISTS "LiteLLM_RetiredAgentIdentity_provider_tenant_id_client_id_key" ON "LiteLLM_RetiredAgentIdentity"("provider", "tenant_id", "client_id");
-- CreateIndex
CREATE INDEX IF NOT EXISTS "LiteLLM_VerifiedSubject_user_id_idx" ON "LiteLLM_VerifiedSubject"("user_id");
-- CreateIndex
CREATE UNIQUE INDEX IF NOT EXISTS "LiteLLM_VerifiedSubject_issuer_tenant_id_oid_key" ON "LiteLLM_VerifiedSubject"("issuer", "tenant_id", "oid");
-- AddForeignKey
DO $$
BEGIN
IF NOT EXISTS (SELECT 1 FROM pg_constraint WHERE conname = 'LiteLLM_AgentIdentity_agent_id_fkey') THEN
ALTER TABLE "LiteLLM_AgentIdentity" ADD CONSTRAINT "LiteLLM_AgentIdentity_agent_id_fkey" FOREIGN KEY ("agent_id") REFERENCES "LiteLLM_AgentsTable"("agent_id") ON DELETE CASCADE ON UPDATE CASCADE;
END IF;
END $$;
-- AddForeignKey
DO $$
BEGIN
IF NOT EXISTS (SELECT 1 FROM pg_constraint WHERE conname = 'LiteLLM_RetiredAgentIdentity_agent_id_fkey') THEN
ALTER TABLE "LiteLLM_RetiredAgentIdentity" ADD CONSTRAINT "LiteLLM_RetiredAgentIdentity_agent_id_fkey" FOREIGN KEY ("agent_id") REFERENCES "LiteLLM_AgentsTable"("agent_id") ON DELETE SET NULL ON UPDATE CASCADE;
END IF;
END $$;
-- AddForeignKey
DO $$
BEGIN
IF NOT EXISTS (SELECT 1 FROM pg_constraint WHERE conname = 'LiteLLM_VerifiedSubject_user_id_fkey') THEN
ALTER TABLE "LiteLLM_VerifiedSubject" ADD CONSTRAINT "LiteLLM_VerifiedSubject_user_id_fkey" FOREIGN KEY ("user_id") REFERENCES "LiteLLM_UserTable"("user_id") ON DELETE CASCADE ON UPDATE CASCADE;
END IF;
END $$;

View file

@ -0,0 +1,2 @@
-- AlterTable
ALTER TABLE "LiteLLM_MCPServerTable" ADD COLUMN IF NOT EXISTS "pinned_tools" JSONB DEFAULT '{}';

View file

@ -0,0 +1,19 @@
CREATE TABLE IF NOT EXISTS "LiteLLM_DailyModelUsage" (
"date" TEXT NOT NULL,
"model_group" TEXT NOT NULL,
"model" TEXT NOT NULL,
"custom_llm_provider" TEXT NOT NULL,
"task_type" TEXT NOT NULL,
"spend" DOUBLE PRECISION NOT NULL DEFAULT 0.0,
"prompt_tokens" BIGINT NOT NULL DEFAULT 0,
"completion_tokens" BIGINT NOT NULL DEFAULT 0,
"request_count" BIGINT NOT NULL DEFAULT 0,
"successful_requests" BIGINT NOT NULL DEFAULT 0,
"failed_requests" BIGINT NOT NULL DEFAULT 0,
"created_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"updated_at" TIMESTAMP(3) NOT NULL,
CONSTRAINT "LiteLLM_DailyModelUsage_pkey" PRIMARY KEY ("date", "model_group", "model", "custom_llm_provider", "task_type")
);
CREATE INDEX IF NOT EXISTS "LiteLLM_DailyModelUsage_date_idx" ON "LiteLLM_DailyModelUsage"("date");
CREATE INDEX IF NOT EXISTS "LiteLLM_DailyModelUsage_model_group_idx" ON "LiteLLM_DailyModelUsage"("model_group");

View file

@ -0,0 +1,10 @@
CREATE TABLE IF NOT EXISTS "LiteLLM_Engine" (
"id" TEXT NOT NULL PRIMARY KEY,
"version" INTEGER NOT NULL DEFAULT 0,
"data" JSONB NOT NULL
);
CREATE TABLE IF NOT EXISTS "LiteLLM_EngineWorker" (
"id" TEXT NOT NULL PRIMARY KEY,
"token_hash" TEXT NOT NULL UNIQUE,
"data" JSONB NOT NULL
);

View file

@ -0,0 +1,7 @@
CREATE TABLE IF NOT EXISTS "LiteLLM_EngineRun" (
"id" TEXT NOT NULL PRIMARY KEY,
"engine_id" TEXT NOT NULL,
"created_at" TIMESTAMP(3) NOT NULL,
"data" JSONB NOT NULL
);
CREATE INDEX IF NOT EXISTS "LiteLLM_EngineRun_engine_id_created_at_idx" ON "LiteLLM_EngineRun"("engine_id", "created_at");

View file

@ -78,6 +78,11 @@ model LiteLLM_AgentsTable {
object_permission_id String?
object_permission LiteLLM_ObjectPermissionTable? @relation(fields: [object_permission_id], references: [object_permission_id])
spend Float @default(0.0)
identity_managed Boolean @default(false)
enabled Boolean @default(true)
execution_mode String @default("autonomous")
identity LiteLLM_AgentIdentity?
retired_identities LiteLLM_RetiredAgentIdentity[]
tpm_limit Int?
rpm_limit Int?
session_tpm_limit Int?
@ -88,6 +93,56 @@ model LiteLLM_AgentsTable {
updated_by String
}
model LiteLLM_AgentIdentity {
agent_id String @id
active Boolean @default(true)
agent LiteLLM_AgentsTable @relation(fields: [agent_id], references: [agent_id], onDelete: Cascade)
provider String
issuer String
tenant_id String
client_id String
service_principal_id String?
required_roles String[] @default([])
required_scopes String[] @default(["user_impersonation"])
revision String @default(uuid())
last_authenticated_at DateTime?
@@unique([provider, tenant_id, client_id])
@@unique([issuer, service_principal_id])
}
model LiteLLM_RetiredAgentIdentity {
binding_id String @id @default(uuid())
agent_id String?
agent LiteLLM_AgentsTable? @relation(fields: [agent_id], references: [agent_id], onDelete: SetNull)
provider String
issuer String
tenant_id String
client_id String
@@unique([provider, tenant_id, client_id])
}
model LiteLLM_RetiredAgent {
original_agent_id String @id
retired_at DateTime @default(now())
}
model LiteLLM_VerifiedSubject {
subject_id String @id @default(uuid())
issuer String
tenant_id String
oid String
kind String @default("human")
user_id String?
user LiteLLM_UserTable? @relation(fields: [user_id], references: [user_id], onDelete: Cascade)
verified_via String @default("sso_interactive")
verified_at DateTime @default(now())
@@unique([issuer, tenant_id, oid])
@@index([user_id])
}
model LiteLLM_OrganizationTable {
organization_id String @id @default(uuid())
organization_alias String
@ -241,6 +296,7 @@ model LiteLLM_DeletedTeamTable {
// Track spend, rate limit, budget Users
model LiteLLM_UserTable {
verified_subjects LiteLLM_VerifiedSubject[]
user_id String @id
user_alias String?
team_id String?
@ -322,6 +378,7 @@ model LiteLLM_MCPServerTable {
allowed_tools String[] @default([])
tool_name_to_display_name Json? @default("{}")
tool_name_to_description Json? @default("{}")
pinned_tools Json? @default("{}")
extra_headers String[] @default([])
static_headers Json? @default("{}")
// Admin-configured environment variables interpolated into static_headers
@ -674,6 +731,7 @@ model LiteLLM_SpendLogs {
session_id String?
status String?
mcp_namespaced_tool_name String?
billing_agent_id String?
agent_id String?
proxy_server_request Json? @default("{}")
litellm_call_id String?
@ -1259,6 +1317,26 @@ model LiteLLM_DailyToolSpend {
@@id([date, tool_name])
}
model LiteLLM_DailyModelUsage {
date String
model_group String
model String
custom_llm_provider String
task_type String
spend Float @default(0.0)
prompt_tokens BigInt @default(0)
completion_tokens BigInt @default(0)
request_count BigInt @default(0)
successful_requests BigInt @default(0)
failed_requests BigInt @default(0)
created_at DateTime @default(now())
updated_at DateTime @updatedAt
@@id([date, model_group, model, custom_llm_provider, task_type])
@@index([date])
@@index([model_group])
}
// Gateway request counts recorded at the ASGI edge by
// BillableRequestMetricsMiddleware. This is the source of truth for SGR
// (successful gateway requests): it counts what the proxy actually answered,
@ -1816,3 +1894,24 @@ model LiteLLM_WorkflowMessage {
@@unique([run_id, sequence_number])
@@index([run_id])
}
model LiteLLM_Engine {
id String @id
version Int @default(0)
data Json
}
model LiteLLM_EngineRun {
id String @id
engine_id String
created_at DateTime
data Json
@@index([engine_id, created_at])
}
model LiteLLM_EngineWorker {
id String @id
token_hash String @unique
data Json
}

View file

@ -1,6 +1,6 @@
[project]
name = "litellm-proxy-extras"
version = "0.4.102"
version = "0.4.103"
description = "Additional files for the LiteLLM Proxy. Reduces the size of the main litellm package."
readme = "README.md"
requires-python = ">=3.9"
@ -26,7 +26,7 @@ required-version = ">=0.10.9"
module-root = ""
[tool.commitizen]
version = "0.4.102"
version = "0.4.103"
version_files = [
"pyproject.toml:^version",
"../pyproject.toml:litellm-proxy-extras==",

View file

@ -0,0 +1,22 @@
---
name: rust-tracing
description: Add or change Rust diagnostic tracing in litellm-rust, including route spans, subscriber layers, and Python logger delivery
---
# Rust tracing
Use upstream `tracing` throughout Rust, including `#[tracing::instrument]`, events, and span propagation. Centralize collection and delivery infrastructure in `crates/tracing`. Direct upstream imports still reach our configured subscriber; re-exporting macros does not control delivery. Do not introduce Rust `log` or `pyo3-log` for this path
`litellm-tracing` owns shared subscriber layers, span field collection, and diagnostic processing. Keep adapters composable as `tracing_subscriber::Layer`s, with `Logger` providing host setup. Runtime-specific delivery belongs in the host bridge. The Python bridge delivers directly to the existing Python SDK logger, preserving its handlers, filtering, redaction, and request correlation. Keep Python dependencies out of `crates/tracing`
Hosts configure subscribers. Keep Python execution scoped to its captured dispatch rather than installing a process-wide subscriber. Propagate both span context and dispatch across spawned work and returned streams
In core, instrument execution shared by native calls and hosted machines. Use consistent route, model, provider, streaming, and outcome fields. Put status recording at shared provider boundaries instead of scattering basic logging through handlers. Keep upstream HTTP status separate from route success
Use `skip_all` and explicitly selected fields. Basic tracing excludes bodies, credentials, headers, and raw error strings. Avoid automatic `ret` or `err` capture of sensitive values. Keep payload diagnostics separate and subject to existing redaction
A returned stream retains its route span until exhaustion, error, or drop, with exactly one terminal outcome. Builder construction does not start a trace. Never hold a span entry guard across an await. Diagnostic tracing remains separate from lifecycle callbacks and `CustomLogger` dispatch
Use `litellm_tracing::sink_layer` to compose a sink with other subscriber layers. It inherits span fields into events and emits span-close summaries with elapsed time. Test observable records, concurrent isolation, dynamic filtering, sensitive-field exclusion, and stream cancellation when changing this behavior
Consult the [tracing API](https://docs.rs/tracing/latest/tracing/) and [subscriber layers](https://docs.rs/tracing-subscriber/latest/tracing_subscriber/layer/index.html) for implementation details

View file

@ -1,5 +1,7 @@
# Rust workspace rules
For diagnostic tracing changes, follow [.agents/skills/rust-tracing/SKILL.md](.agents/skills/rust-tracing/SKILL.md)
## Test placement
- Never create a `tests.rs` (or `test.rs`) file under `src/`, and never `#[path = "tests.rs"] mod tests;`
@ -16,7 +18,9 @@ Use [`#[rstest]`](https://docs.rs/rstest/latest/rstest/attr.rstest.html) for new
## Error definitions
- A crate's errors live in `src/error.rs`, defined with `thiserror`, and re-exported from `lib.rs`
- Put message templates in the variant's `#[error(...)]` declaration. Callers pass only the small typed arguments needed to fill them, never `Error::Variant(format!(...))` or a preformatted message. Keep the smallest set of neutral variants that callers need to distinguish; different wording or providers do not justify new variants
- Default to one top-level `Error` enum per crate, with one variant per failure mode and a `#[error(...)]` message on each. A failure mode is something a caller handles differently (phase, status code, retry, a message Python parity pins exactly); failures no caller tells apart share one variant and differ only in its message
- Keep shared error enums minimal and provider-neutral. Provider names, credential types, configuration fields, and setup guidance belong in caller-supplied data, not dedicated variants or hardcoded shared messages. Reuse a variant for the same failure mode across providers, such as `MissingApiBase { provider: "Azure", guidance: "..." }`. An exact parity message does not justify a provider-specific variant when caller-supplied context can preserve it
- Wrap a lower-level error as a variant with `#[from]` or `#[source]` instead of flattening it to a string
- Exception: split into separate types when different functions fail in disjoint ways, especially when different callers see them. A shared enum would force every caller to match variants its function can never return
- Name a split type after what went wrong (a unit struct is fine for a single failure mode), not after the function that returns it

1552
litellm-rust/Cargo.lock generated

File diff suppressed because it is too large Load diff

View file

@ -12,12 +12,18 @@ repository = "https://github.com/BerriAI/litellm"
litellm-config = { path = "crates/config" }
litellm-router = { path = "crates/router" }
litellm-tracing = { path = "crates/tracing" }
litellm-traces = { path = "crates/traces" }
litellm-core = { path = "crates/core" }
litellm-gateway-mcp = { path = "crates/gateway-mcp" }
litellm-gateway = { path = "crates/gateway" }
litellm-gateway-inference = { path = "crates/gateway-inference" }
litellm-gateway-auth = { path = "crates/gateway-auth" }
litellm-gateway-management = { path = "crates/gateway-management" }
litellm-gateway-ui = { path = "crates/gateway-ui" }
litellm-coroutine = { path = "crates/coroutine" }
litellm-host = { path = "crates/host" }
litellm-host-http = { path = "crates/host-http" }
litellm-host-native = { path = "crates/host-native" }
litellm-callbacks-legacy-python = { path = "crates/callbacks-legacy-python" }
litellm-framing = { path = "crates/framer" }
litellm-auth = { path = "crates/auth" }
@ -34,8 +40,10 @@ litellm-secrets-azure = { path = "crates/secrets-azure" }
litellm-secrets-cyberark = { path = "crates/secrets-cyberark" }
litellm-http = { path = "crates/http" }
litellm-llms = { path = "crates/llms" }
litellm-types = { path = "crates/types" }
litellm-llms-types = { path = "crates/llms-types" }
litellm-core-utils = { path = "crates/core-utils" }
litellm-db = { path = "crates/db" }
litellm-db-testing = { path = "crates/db-testing" }
litellm-cache = { path = "crates/cache" }
litellm-cache-azure-blob = { path = "crates/cache-azure-blob" }
litellm-cache-memory = { path = "crates/cache-memory" }
@ -56,6 +64,8 @@ litellm-python-compat = { path = "crates/python-compat" }
tracing = "0.1"
axum = { version = "0.8.9", default-features = false, features = ["http1", "tokio", "multipart"] }
axum-login = "0.18.0"
tower-sessions = { version = "0.14.0", features = ["memory-store"] }
bytes = "1"
http = "1"
google-cloud-auth = { version = "1.16.0", default-features = false }
@ -65,6 +75,7 @@ proptest = "1.7.0"
pyo3 = "0.29.2"
pyo3-async-runtimes = { version = "0.29.0", features = ["tokio-runtime"] }
rand = "0.8"
macro_rules_attribute = "0.2.3"
schemars = "1"
reqwest = { version = "0.12", default-features = false, features = ["json", "multipart", "rustls-tls", "http2", "stream"] }
qdrant-client = { version = "1.19.0", default-features = false }
@ -80,6 +91,7 @@ serde = { version = "1.0", features = ["derive"] }
serde_json = { version = "1.0", features = ["float_roundtrip"] }
serde_with = { version = "=3.16.1", default-features = false, features = ["std", "macros"] }
sha2 = "0.10"
sqlx = { version = "0.9.0", default-features = false, features = ["json", "macros", "postgres", "runtime-tokio", "chrono", "tls-rustls-ring-native-roots"] }
subtle = "2"
thiserror = "2.0"
tokenizers = { version = "0.23.1", default-features = false, features = ["onig"] }

View file

@ -1,4 +1,4 @@
# The Tokio runtime is reached only through `host-python/src/execution.rs`, whose fork gate
# The Tokio runtime is reached only through `host-python/src/runtime.rs`, whose fork gate
# must see every entry. Going around it makes a fork-after-use hang instead of raising.
disallowed-methods = [
{ path = "pyo3_async_runtimes::tokio::get_runtime", reason = "use litellm_host_python::run_sync / run_sync_value" },
@ -12,6 +12,13 @@ disallowed-methods = [
{ path = "reqwest::ClientBuilder::danger_accept_invalid_certs", reason = "set HttpClientConfig::verify instead" },
{ path = "reqwest::ClientBuilder::identity", reason = "set HttpClientConfig::client_certificate instead" },
{ path = "reqwest::ClientBuilder::use_preconfigured_tls", reason = "HttpClientConfig owns the TLS configuration" },
{ path = "sqlx::query", reason = "use sqlx::query! or query_file! so the SQL is checked against the migrated schema" },
{ path = "sqlx::query_as", reason = "use sqlx::query_as! or query_file_as! so the SQL is checked against the migrated schema" },
{ path = "sqlx::query_scalar", reason = "use sqlx::query_scalar! so the SQL is checked against the migrated schema" },
{ path = "sqlx::query_with", reason = "use sqlx::query! or query_file! so the SQL is checked against the migrated schema" },
{ path = "sqlx::query_as_with", reason = "use sqlx::query_as! or query_file_as! so the SQL is checked against the migrated schema" },
{ path = "sqlx::query_scalar_with", reason = "use sqlx::query_scalar! so the SQL is checked against the migrated schema" },
{ path = "sqlx::raw_sql", reason = "raw_sql is unchecked; use the checked query macros" },
]
# Every outbound client comes from litellm_http::HttpClientPool so it honors the host's TLS,

View file

@ -133,7 +133,7 @@ impl NativeAzureTokenAcquirer {
let token = credential
.get_token(&[scope.as_str()], None)
.await
.map_err(|error| Error::AzureTokenAcquisition(error.to_string()))?;
.map_err(|error| Error::CredentialAcquisition(error.to_string().into()))?;
let expires_on = u64::try_from(token.expires_on.unix_timestamp())
.ok()
.map(|seconds| UNIX_EPOCH + Duration::from_secs(seconds));
@ -250,7 +250,12 @@ fn validate_authority(request: &NativeAzureRequest) -> Result<(), Error> {
let Some(authority) = authority else {
return Ok(());
};
let url = url::Url::parse(authority.value()).map_err(|_| Error::InvalidAzureAuthority)?;
let url = url::Url::parse(authority.value()).map_err(|_| {
Error::InvalidConfiguration(
"Azure authority must be an HTTPS origin without credentials, query, or fragment"
.into(),
)
})?;
if url.scheme() != "https"
|| url.host_str().is_none()
|| !url.username().is_empty()
@ -259,7 +264,10 @@ fn validate_authority(request: &NativeAzureRequest) -> Result<(), Error> {
|| url.fragment().is_some()
|| !matches!(url.path(), "" | "/")
{
return Err(Error::InvalidAzureAuthority);
return Err(Error::InvalidConfiguration(
"Azure authority must be an HTTPS origin without credentials, query, or fragment"
.into(),
));
}
Ok(())
}
@ -368,7 +376,9 @@ fn trusted_source(sources: &[InputSource]) -> InputSource {
}
fn mixed_sources<T>() -> Result<T, Error> {
Err(Error::MixedAzureCredentialSources)
Err(Error::InvalidConfiguration(
"request-controlled Azure auth inputs cannot be combined with host credentials".into(),
))
}
fn build_credential(
@ -433,7 +443,12 @@ fn build_credential(
NativeAzureRequest::DeveloperTools { .. } => DeveloperToolsCredential::new(None)
.map(|credential| credential as Arc<dyn TokenCredential>),
}
.map_err(|error| Error::AzureCredentialInitialization(error.to_string()))
.map_err(|error| {
Error::InvalidConfiguration(litellm_auth_types::ErrorDetail::failed(
"Azure credential initialization",
error,
))
})
}
fn client_options(
@ -638,7 +653,7 @@ mod tests {
assert_eq!(transport.requests.lock().unwrap().len(), 6);
}
#[test]
#[rstest::rstest]
fn request_authority_requires_request_owned_client_secret_identity() {
let error = ValidatedAzureRequest::new(sourced_client_secret(
InputSource::Deployment,
@ -647,10 +662,13 @@ mod tests {
))
.unwrap_err();
assert!(matches!(
assert_eq!(
error,
litellm_auth_types::Error::MixedAzureCredentialSources
));
litellm_auth_types::Error::InvalidConfiguration(
"request-controlled Azure auth inputs cannot be combined with host credentials"
.into()
)
);
}
#[test]
@ -665,24 +683,24 @@ mod tests {
assert_eq!(request.credential_source(), InputSource::Request);
}
#[test]
fn authority_is_restricted_to_an_https_origin() {
for authority in [
"http://login.example",
"https://user@login.example",
"https://login.example/tenant",
"https://login.example?target=other",
] {
let error = ValidatedAzureRequest::new(sourced_client_secret(
InputSource::Deployment,
InputSource::Deployment,
authority,
))
.unwrap_err();
assert!(matches!(
error,
litellm_auth_types::Error::InvalidAzureAuthority
));
}
#[rstest::rstest]
#[case::http("http://login.example")]
#[case::userinfo("https://user@login.example")]
#[case::path("https://login.example/tenant")]
#[case::query("https://login.example?target=other")]
fn authority_is_restricted_to_an_https_origin(#[case] authority: &str) {
let error = ValidatedAzureRequest::new(sourced_client_secret(
InputSource::Deployment,
InputSource::Deployment,
authority,
))
.unwrap_err();
assert_eq!(
error,
litellm_auth_types::Error::InvalidConfiguration(
"Azure authority must be an HTTPS origin without credentials, query, or fragment"
.into()
)
);
}
}

View file

@ -91,7 +91,9 @@ impl AzureAuthService {
AzureCredentialPlan::Caller(caller) => {
let credential = caller.acquire().await?;
if credential.secret().expose().is_empty() {
return Err(Error::EmptyAzureToken);
return Err(Error::EmptyCallerCredential(
"Azure AD token provider returned an empty token",
));
}
Ok(Some(Sourced::new(credential, InputSource::Deployment)))
}
@ -104,7 +106,11 @@ impl AzureAuthService {
} => {
let assertion = resolve_reference(inputs, env_lookup, reference.value())
.await?
.ok_or(Error::UnresolvedOidcReference)?;
.ok_or_else(|| {
Error::CredentialAcquisition(
"Azure OIDC reference did not resolve to a value".into(),
)
})?;
let request = ValidatedAzureRequest::new(NativeAzureRequest::ClientAssertion {
tenant_id,
client_id,
@ -167,7 +173,7 @@ pub(crate) fn select_auth_plan(
.map(|selector| Sourced::new(selector, value.source()))
})
.transpose()
.map_err(|_| Error::InvalidAzureSelector)?;
.map_err(|_| Error::InvalidConfiguration("invalid Azure credential selector".into()))?;
let federated_token_file = configured_string(
&inputs.federated_token_file,
AZURE_FEDERATED_TOKEN_FILE_ENV,
@ -257,7 +263,9 @@ fn select_native_plan(
let selection_source = selected.source();
match selected.into_value() {
AzureCredentialType::ClientSecretCredential => Err(Error::MissingClientSecretFields),
AzureCredentialType::ClientSecretCredential => Err(Error::InvalidConfiguration(
"ClientSecretCredential requires tenant_id, client_id, and client_secret".into(),
)),
AzureCredentialType::WorkloadIdentityCredential => {
Ok(AzureCredentialPlan::Native(ValidatedAzureRequest::new(
workload_request(tenant_id, client_id, federated_token_file, scope, authority)?,
@ -341,9 +349,17 @@ fn workload_request(
authority: Option<Sourced<String>>,
) -> Result<NativeAzureRequest, Error> {
Ok(NativeAzureRequest::WorkloadIdentity {
tenant_id: tenant_id.ok_or(Error::MissingWorkloadTenant)?,
client_id: client_id.ok_or(Error::MissingWorkloadClient)?,
token_file_path: token_file_path.ok_or(Error::MissingWorkloadTokenFile)?,
tenant_id: tenant_id.ok_or_else(|| {
Error::InvalidConfiguration("WorkloadIdentityCredential requires tenant_id".into())
})?,
client_id: client_id.ok_or_else(|| {
Error::InvalidConfiguration("WorkloadIdentityCredential requires client_id".into())
})?,
token_file_path: token_file_path.ok_or_else(|| {
Error::InvalidConfiguration(
"WorkloadIdentityCredential requires azure_federated_token_file".into(),
)
})?,
scope,
authority,
})
@ -394,10 +410,11 @@ async fn resolve_reference(
.map_or(CredentialLookup::Missing, CredentialLookup::Found),
CredentialRef::None => return Ok(None),
CredentialRef::File(_) | CredentialRef::Request(_) | CredentialRef::Host(_) => {
let resolver = inputs
.credential_resolver
.as_ref()
.ok_or(Error::MissingHostResolver)?;
let resolver = inputs.credential_resolver.as_ref().ok_or_else(|| {
Error::InvalidConfiguration(
"credential reference requires a host credential resolver".into(),
)
})?;
resolver.resolve(reference).await?
}
};
@ -415,7 +432,9 @@ fn oidc_reference(
};
let value = token.value().expose();
if token.source() == InputSource::Request && value.starts_with("oidc/") {
return Err(Error::RequestAzureCredentialReference);
return Err(Error::InvalidConfiguration(
"request-controlled Azure credential references are not allowed".into(),
));
}
if let Some(name) = value.strip_prefix("oidc/env/") {
return non_empty_reference(name, "OIDC environment reference")
@ -437,14 +456,20 @@ fn oidc_reference(
)));
}
if value.starts_with("oidc/") {
return Err(Error::UnsupportedOidcReference);
return Err(Error::InvalidConfiguration(
"unsupported OIDC reference".into(),
));
}
Ok(None)
}
fn non_empty_reference(value: &str, kind: &str) -> Result<String, Error> {
if value.is_empty() {
return Err(Error::EmptyReference(kind.to_string()));
return Err(Error::InvalidConfiguration(
litellm_auth_types::ErrorDetail::Empty {
subject: kind.into(),
},
));
}
Ok(value.to_string())
}
@ -493,7 +518,7 @@ mod tests {
expires_on: None,
})
} else {
Err(Error::AzureTokenAcquisition(format!("{kind} failed")))
Err(Error::CredentialAcquisition(kind.into()))
}
})
}
@ -602,7 +627,7 @@ mod tests {
assert!(error.to_string().contains("unsupported OIDC reference"));
}
#[test]
#[rstest::rstest]
fn request_oidc_reference_is_rejected_before_lookup() {
let params = json!({
"azure_ad_token": "oidc/env/ASSERTION",
@ -624,7 +649,12 @@ mod tests {
})
.unwrap_err();
assert!(matches!(error, Error::RequestAzureCredentialReference));
assert_eq!(
error,
Error::InvalidConfiguration(
"request-controlled Azure credential references are not allowed".into()
)
);
}
#[tokio::test]
@ -723,6 +753,7 @@ mod tests {
assert_eq!(credential.value().secret().expose(), "caller-token");
}
#[rstest::rstest]
#[tokio::test]
async fn empty_caller_token_is_rejected() {
let error = AzureAuthService::default()
@ -730,6 +761,9 @@ mod tests {
.await
.unwrap_err();
assert!(matches!(error, Error::EmptyAzureToken));
assert_eq!(
error,
Error::EmptyCallerCredential("Azure AD token provider returned an empty token")
);
}
}

View file

@ -117,7 +117,12 @@ fn string_config(
None => Ok(ConfigValue::Absent),
Some(Value::Null) => Ok(ConfigValue::ExplicitNone(source)),
Some(Value::String(value)) => Ok(ConfigValue::Value(Sourced::new(value.clone(), source))),
Some(_) => Err(Error::InvalidFieldType(name.to_string())),
Some(_) => Err(Error::InvalidConfiguration(
litellm_auth_types::ErrorDetail::InvalidType {
field: name.into(),
expected: "a string or null",
},
)),
}
}

View file

@ -19,3 +19,6 @@ tokio.workspace = true
gcp_auth = "0.12.7"
google-cloud-auth = { workspace = true, optional = true }
http = { workspace = true, optional = true }
[dev-dependencies]
rstest.workspace = true

View file

@ -299,7 +299,7 @@ fn validate_request_credentials(configured: &str) -> Result<&str, Error> {
.map(str::to_string)
});
if token_uri.as_deref() != Some(GOOGLE_OAUTH_TOKEN_ENDPOINT) {
return Err(Error::RequestVertexTokenEndpoint);
return Err(Error::InvalidConfiguration("request-controlled Vertex credentials must use the canonical Google OAuth token endpoint".into()));
}
Ok(configured)
}
@ -376,10 +376,20 @@ fn optional_credentials(
.map(SecretValue::new)
.map(|value| Sourced::new(value, source))
.map(Some)
.map_err(|error| Error::InvalidFieldType(format!("{}: {error}", names[0])));
.map_err(|error| {
Error::InvalidConfiguration(litellm_auth_types::ErrorDetail::failed(
"credential serialization",
error,
))
});
}
Some(_) => {
return Err(Error::InvalidFieldType(names[0].to_string()));
return Err(Error::InvalidConfiguration(
litellm_auth_types::ErrorDetail::InvalidType {
field: names[0].into(),
expected: "a string or null",
},
));
}
}
}
@ -397,7 +407,12 @@ fn optional_string(params: &Map<String, Value>, names: &[&str]) -> Result<Option
Some(Value::String(value)) if value.trim().is_empty() => continue,
Some(Value::String(value)) => return Ok(Some(value.clone())),
Some(_) => {
return Err(Error::InvalidFieldType(names[0].to_string()));
return Err(Error::InvalidConfiguration(
litellm_auth_types::ErrorDetail::InvalidType {
field: names[0].into(),
expected: "a string or null",
},
));
}
}
}
@ -411,7 +426,10 @@ fn non_empty_env(env_lookup: &dyn Fn(&str) -> Option<String>, name: &str) -> Opt
}
fn auth_acquisition_error(error: gcp_auth::Error) -> Error {
Error::VertexTokenAcquisition(error.to_string())
Error::CredentialAcquisition(litellm_auth_types::ErrorDetail::failed(
"Vertex AI credentials",
error,
))
}
#[cfg(test)]
@ -612,20 +630,15 @@ mod tests {
);
}
#[test]
fn request_credentials_require_canonical_token_endpoint() {
assert!(
validate_request_credentials(r#"{"token_uri":"https://oauth2.googleapis.com/token"}"#)
.is_ok()
);
assert!(matches!(
validate_request_credentials(r#"{"token_uri":"http://127.0.0.1/token"}"#),
Err(Error::RequestVertexTokenEndpoint)
));
assert!(matches!(
validate_request_credentials("{}"),
Err(Error::RequestVertexTokenEndpoint)
));
#[rstest::rstest]
#[case::canonical_endpoint(r#"{"token_uri":"https://oauth2.googleapis.com/token"}"#, true)]
#[case::noncanonical_endpoint(r#"{"token_uri":"http://127.0.0.1/token"}"#, false)]
#[case::missing_endpoint("{}", false)]
fn request_credentials_require_canonical_token_endpoint(
#[case] credentials: &str,
#[case] accepted: bool,
) {
assert_eq!(validate_request_credentials(credentials).is_ok(), accepted);
}
#[tokio::test]

View file

@ -12,4 +12,5 @@ thiserror.workspace = true
veil.workspace = true
[dev-dependencies]
rstest.workspace = true
tokio.workspace = true

View file

@ -86,7 +86,9 @@ impl CredentialPlan {
Self::Caller(caller) => {
let credential = caller.acquire().await?;
if credential.secret().expose().is_empty() {
return Err(Error::EmptyCallerCredential);
return Err(Error::EmptyCallerCredential(
"credential caller returned an empty credential",
));
}
Ok(CredentialPlanResolution::Resolved(credential))
}
@ -147,10 +149,15 @@ mod tests {
impl CredentialResolver for FailingResolver {
fn resolve<'a>(&'a self, _reference: &'a CredentialRef) -> CredentialLookupFuture<'a> {
Box::pin(async { Err(Error::UnresolvedOidcReference) })
Box::pin(async {
Err(Error::CredentialAcquisition(
"host credential lookup failed".into(),
))
})
}
}
#[rstest::rstest]
#[tokio::test]
async fn acquisition_failure_is_terminal() {
let resolver = CredentialResolverHandle::new(Arc::new(FailingResolver));
@ -161,6 +168,9 @@ mod tests {
.await
.expect_err("acquisition errors cannot become fallback");
assert_eq!(error, Error::UnresolvedOidcReference);
assert_eq!(
error,
Error::CredentialAcquisition("host credential lookup failed".into())
);
}
}

View file

@ -2,84 +2,16 @@ use thiserror::Error as ThisError;
#[derive(Clone, Debug, ThisError, PartialEq, Eq)]
pub enum Error {
#[error("invalid authentication configuration: credential header already exists")]
ExistingCredentialHeader,
#[error(
"invalid authentication configuration: credential plan is not allowed by the provider auth policy"
)]
DisallowedCredentialPlan,
#[error("invalid authentication configuration: credential cannot be empty")]
EmptyCredential,
#[error("invalid authentication configuration: invalid Azure credential selector")]
InvalidAzureSelector,
#[error(
"invalid authentication configuration: ClientSecretCredential requires tenant_id, client_id, and client_secret"
)]
MissingClientSecretFields,
#[error("invalid authentication configuration: WorkloadIdentityCredential requires tenant_id")]
MissingWorkloadTenant,
#[error("invalid authentication configuration: WorkloadIdentityCredential requires client_id")]
MissingWorkloadClient,
#[error(
"invalid authentication configuration: WorkloadIdentityCredential requires azure_federated_token_file"
)]
MissingWorkloadTokenFile,
#[error(
"invalid authentication configuration: credential reference requires a host credential resolver"
)]
MissingHostResolver,
#[error(
"invalid authentication configuration: caller credential plan requires provider-specific inputs"
)]
MissingCallerInputs,
#[error("invalid authentication configuration: credential header {0} already exists")]
DuplicateHeader(&'static str),
#[error("invalid authentication configuration: {0} must be a string or null")]
InvalidFieldType(String),
#[error("invalid authentication configuration: unsupported OIDC reference")]
UnsupportedOidcReference,
#[error("invalid authentication configuration: {0} cannot be empty")]
EmptyReference(String),
#[error("invalid authentication configuration: Azure credential initialization failed: {0}")]
AzureCredentialInitialization(String),
#[error(
"invalid authentication configuration: Azure authority must be an HTTPS origin without credentials, query, or fragment"
)]
InvalidAzureAuthority,
#[error(
"invalid authentication configuration: request-controlled Azure auth inputs cannot be combined with host credentials"
)]
MixedAzureCredentialSources,
#[error(
"invalid authentication configuration: request-controlled Azure credential references are not allowed"
)]
RequestAzureCredentialReference,
#[error(
"invalid authentication configuration: host credentials cannot be sent to a request-controlled Azure endpoint"
)]
RequestAzureCredentialDestination,
#[error(
"invalid authentication configuration: credentials cannot be sent to a request-controlled Vertex AI endpoint"
)]
RequestVertexCredentialDestination,
#[error(
"invalid authentication configuration: request-controlled Vertex credentials must use the canonical Google OAuth token endpoint"
)]
RequestVertexTokenEndpoint,
#[error("invalid authentication configuration: {0}")]
InvalidConfiguration(#[source] ErrorDetail),
#[error("credential acquisition failed: {0}")]
AzureTokenAcquisition(String),
#[error("credential acquisition failed: Vertex AI credentials: {0}")]
VertexTokenAcquisition(String),
CredentialAcquisition(#[source] ErrorDetail),
#[error("credential caller failed: {0}")]
EmptyCallerCredential(&'static str),
#[error("{0}")]
ProviderAuthentication(String),
#[error("credential acquisition failed: {}", .0.iter().map(ToString::to_string).collect::<Vec<_>>().join("; "))]
CredentialChain(Vec<Error>),
#[error("credential caller failed: credential caller returned an empty credential")]
EmptyCallerCredential,
#[error("credential caller failed: Azure AD token provider returned an empty token")]
EmptyAzureToken,
#[error("credential acquisition failed: Azure OIDC reference did not resolve to a value")]
UnresolvedOidcReference,
#[error(
"Missing {provider} API Key - Set `api_key` or the {environment_variable} environment variable"
)]
@ -87,34 +19,87 @@ pub enum Error {
provider: &'static str,
environment_variable: &'static str,
},
#[error(
"Missing {provider} API Base - Set {environment_variable} environment variable or pass api_base parameter"
)]
#[error("Missing {provider} API Base - {guidance}")]
MissingApiBase {
provider: &'static str,
environment_variable: &'static str,
guidance: &'static str,
},
#[error(
"Missing Azure API Base - Set `api_base` or the AZURE_API_BASE environment variable. Expected format: https://<resource-name>.services.ai.azure.com/anthropic"
)]
MissingAzureApiBase,
#[error("invalid authentication header")]
InvalidHeader,
}
#[cfg(test)]
mod tests {
use super::Error;
#[derive(Clone, Debug, PartialEq, Eq, thiserror::Error)]
pub enum ErrorDetail {
#[error("{0}")]
Message(String),
#[error("{field} must be {expected}")]
InvalidType {
field: String,
expected: &'static str,
},
#[error("{subject} cannot be empty")]
Empty { subject: String },
#[error("credential header {0} already exists")]
DuplicateHeader(&'static str),
#[error("{operation} failed: {source}")]
Failed {
operation: &'static str,
#[source]
source: ErrorSource,
},
}
#[test]
fn missing_api_key_names_provider_and_environment_variable() {
assert_eq!(
Error::MissingApiKey {
provider: "Anthropic",
environment_variable: "ANTHROPIC_API_KEY",
}
.to_string(),
"Missing Anthropic API Key - Set `api_key` or the ANTHROPIC_API_KEY environment variable"
);
impl ErrorDetail {
pub fn failed(
operation: &'static str,
source: impl std::error::Error + Send + Sync + 'static,
) -> Self {
Self::Failed {
operation,
source: ErrorSource::new(source),
}
}
}
impl From<String> for ErrorDetail {
fn from(message: String) -> Self {
Self::Message(message)
}
}
impl From<&str> for ErrorDetail {
fn from(message: &str) -> Self {
Self::Message(message.into())
}
}
#[derive(Clone, Debug)]
pub struct ErrorSource(std::sync::Arc<dyn std::error::Error + Send + Sync>);
impl std::ops::Deref for ErrorSource {
type Target = dyn std::error::Error + Send + Sync;
fn deref(&self) -> &Self::Target {
self.0.as_ref()
}
}
impl std::fmt::Display for ErrorSource {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
std::fmt::Display::fmt(&self.0, formatter)
}
}
impl ErrorSource {
pub fn new(error: impl std::error::Error + Send + Sync + 'static) -> Self {
Self(std::sync::Arc::new(error))
}
}
impl PartialEq for ErrorSource {
fn eq(&self, other: &Self) -> bool {
std::sync::Arc::ptr_eq(&self.0, &other.0)
}
}
impl Eq for ErrorSource {}

View file

@ -21,13 +21,17 @@ pub fn apply_credential(
placement: CredentialPlacement,
) -> Result<Vec<(String, String)>, Error> {
if credential.trim().is_empty() {
return Err(Error::EmptyCredential);
return Err(Error::InvalidConfiguration(
"credential cannot be empty".into(),
));
}
if headers
.iter()
.any(|(name, _)| name.eq_ignore_ascii_case(placement.header_name()))
{
return Err(Error::DuplicateHeader(placement.header_name()));
return Err(Error::InvalidConfiguration(
crate::ErrorDetail::DuplicateHeader(placement.header_name()),
));
}
let value = match placement {
CredentialPlacement::Bearer => format!("Bearer {credential}"),

View file

@ -50,7 +50,7 @@ pub use credential::{
CredentialFileRef, CredentialLookup, CredentialLookupFuture, CredentialPlan,
CredentialPlanResolution, CredentialRef, CredentialResolver, CredentialResolverHandle,
};
pub use error::Error;
pub use error::{Error, ErrorDetail, ErrorSource};
pub use http::CredentialPlacement;
pub use policy::{CredentialPlanKind, CredentialRule, ExistingHeaderBehavior, ProviderAuthPolicy};
pub use secret::SecretValue;

View file

@ -47,14 +47,20 @@ impl ProviderAuthPolicy {
if self.has_existing_credential(&headers) {
return match self.existing_header_behavior {
ExistingHeaderBehavior::Preserve => Ok(headers),
ExistingHeaderBehavior::Reject => Err(Error::ExistingCredentialHeader),
ExistingHeaderBehavior::Reject => Err(Error::InvalidConfiguration(
"credential header already exists".into(),
)),
};
}
let rule = self
.rules
.iter()
.find(|rule| rule.kind == kind)
.ok_or(Error::DisallowedCredentialPlan)?;
.ok_or_else(|| {
Error::InvalidConfiguration(
"credential plan is not allowed by the provider auth policy".into(),
)
})?;
apply_credential(headers, credential.secret().expose(), rule.placement)
}
}

View file

@ -39,7 +39,33 @@ impl TokenProviderHandle {
Self(caller)
}
pub fn from_callback<F, Fut>(acquire: F) -> Self
where
F: Fn() -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<ResolvedCredential, Error>> + Send + 'static,
{
Self::new(Arc::new(CallbackTokenProvider(acquire)))
}
pub async fn acquire(&self) -> Result<ResolvedCredential, Error> {
self.0.acquire().await
}
}
struct CallbackTokenProvider<F>(F);
impl<F> std::fmt::Debug for CallbackTokenProvider<F> {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("CallbackTokenProvider")
}
}
impl<F, Fut> TokenProvider for CallbackTokenProvider<F>
where
F: Fn() -> Fut + Send + Sync,
Fut: Future<Output = Result<ResolvedCredential, Error>> + Send + 'static,
{
fn acquire(&self) -> TokenFuture<'_> {
Box::pin((self.0)())
}
}

View file

@ -0,0 +1,74 @@
use litellm_auth_types::Error;
use rstest::rstest;
#[rstest]
#[case::api_key(
Error::MissingApiKey { provider: "Example", environment_variable: "EXAMPLE_API_KEY" },
"Missing Example API Key - Set `api_key` or the EXAMPLE_API_KEY environment variable"
)]
#[case::another_api_key(
Error::MissingApiKey { provider: "Custom", environment_variable: "CUSTOM_KEY" },
"Missing Custom API Key - Set `api_key` or the CUSTOM_KEY environment variable"
)]
#[case::api_base(
Error::MissingApiBase { provider: "Example", guidance: "Pass api_base" },
"Missing Example API Base - Pass api_base"
)]
#[case::another_api_base(
Error::MissingApiBase { provider: "Custom", guidance: "Set CUSTOM_ENDPOINT" },
"Missing Custom API Base - Set CUSTOM_ENDPOINT"
)]
#[case::configuration(
Error::InvalidConfiguration("credential selector is invalid".into()),
"invalid authentication configuration: credential selector is invalid"
)]
#[case::acquisition(
Error::CredentialAcquisition("token expired".into()),
"credential acquisition failed: token expired"
)]
#[case::caller(
Error::EmptyCallerCredential("empty token"),
"credential caller failed: empty token"
)]
#[case::provider(
Error::ProviderAuthentication("provider rejected credentials".into()),
"provider rejected credentials"
)]
#[case::chain(
Error::CredentialChain(vec![
Error::CredentialAcquisition("token expired".into()),
Error::EmptyCallerCredential("empty token"),
]),
"credential acquisition failed: credential acquisition failed: token expired; credential caller failed: empty token"
)]
fn display_preserves_failure_phase_and_caller_context(
#[case] error: Error,
#[case] expected: &str,
) {
assert_eq!(error.to_string(), expected);
}
#[rstest]
#[case::configuration(true)]
#[case::acquisition(false)]
fn contextual_failures_keep_the_original_source(#[case] configuration: bool) {
use litellm_auth_types::ErrorDetail;
let detail = ErrorDetail::failed(
"test credential",
std::io::Error::from(std::io::ErrorKind::PermissionDenied),
);
let error = if configuration {
Error::InvalidConfiguration(detail)
} else {
Error::CredentialAcquisition(detail)
};
let source = std::iter::successors(Some(&error as &dyn std::error::Error), |error| {
error.source()
})
.find_map(|error| error.downcast_ref::<std::io::Error>())
.expect("the original credential error remains available");
assert_eq!(source.kind(), std::io::ErrorKind::PermissionDenied);
assert!(error.to_string().contains("test credential failed:"));
assert!(error.to_string().ends_with(&source.to_string()));
}

View file

@ -0,0 +1,104 @@
use std::{
error::Error as StdError,
future::{Future, poll_fn},
sync::{
Arc,
atomic::{AtomicBool, AtomicUsize, Ordering},
},
task::Poll,
time::{Duration, SystemTime},
};
use litellm_auth_types::{
Error, ErrorDetail, ResolvedCredential, SecretValue, TokenProviderHandle,
};
use rstest::rstest;
fn credential(index: usize, access_token: bool) -> ResolvedCredential {
let token = SecretValue::new(format!("credential-{index}"));
if access_token {
return ResolvedCredential::AccessToken {
token,
expires_on: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(index as u64)),
};
}
ResolvedCredential::Static(token)
}
#[rstest]
#[case::static_secret(false)]
#[case::access_token(true)]
#[tokio::test]
async fn callbacks_acquire_fresh_credentials_on_demand(#[case] access_token: bool) {
let calls = Arc::new(AtomicUsize::new(0));
let callback_calls = calls.clone();
let provider = TokenProviderHandle::from_callback(move || {
let index = callback_calls.fetch_add(1, Ordering::SeqCst);
async move {
tokio::task::yield_now().await;
Ok(credential(index, access_token))
}
});
let cloned = provider.clone();
assert_eq!(calls.load(Ordering::SeqCst), 0);
assert_eq!(
provider.acquire().await.unwrap(),
credential(0, access_token)
);
assert_eq!(calls.load(Ordering::SeqCst), 1);
assert_eq!(cloned.acquire().await.unwrap(), credential(1, access_token));
assert_eq!(calls.load(Ordering::SeqCst), 2);
}
#[rstest]
#[tokio::test]
async fn callback_errors_preserve_the_original_source() {
let provider = TokenProviderHandle::from_callback(|| async {
Err(Error::CredentialAcquisition(ErrorDetail::failed(
"caller credential",
std::io::Error::from(std::io::ErrorKind::PermissionDenied),
)))
});
let error = provider.acquire().await.unwrap_err();
assert!(matches!(error, Error::CredentialAcquisition(_)));
let source = std::iter::successors(Some(&error as &(dyn StdError + 'static)), |error| {
(*error).source()
})
.find_map(|error| error.downcast_ref::<std::io::Error>())
.unwrap();
assert_eq!(source.kind(), std::io::ErrorKind::PermissionDenied);
}
struct Release(Arc<AtomicBool>);
impl Drop for Release {
fn drop(&mut self) {
self.0.store(true, Ordering::SeqCst);
}
}
#[rstest]
#[tokio::test]
async fn cancelling_acquisition_drops_the_callback_future() {
let released = Arc::new(AtomicBool::new(false));
let callback_released = released.clone();
let provider = TokenProviderHandle::from_callback(move || {
let released = callback_released.clone();
async move {
let _release = Release(released);
std::future::pending().await
}
});
let mut acquisition = Box::pin(provider.acquire());
poll_fn(|context| {
assert!(acquisition.as_mut().poll(context).is_pending());
assert!(!released.load(Ordering::SeqCst));
Poll::Ready(())
})
.await;
drop(acquisition);
assert!(released.load(Ordering::SeqCst));
}

View file

@ -9,7 +9,7 @@ repository.workspace = true
litellm-cache.workspace = true
py_literal = "0.4.0"
rand.workspace = true
rusqlite = { version = "0.40", features = ["bundled"] }
rusqlite = { version = "0.39", features = ["bundled"] }
serde-pickle = "1.2"
serde_json.workspace = true
tokio.workspace = true

View file

@ -56,6 +56,7 @@ async fn set_writes_encoded_object_and_headers(#[future(awt)] server: MockServer
)]
#[case::missing("missing", ResponseTemplate::new(404), Ok(None))]
#[case::server_error("server-error", ResponseTemplate::new(500), Err(Error::Unavailable))]
#[case::unauthorized("unauthorized", ResponseTemplate::new(401), Err(Error::Unavailable))]
#[case::invalid(
"invalid",
ResponseTemplate::new(200).set_body_string("not json"),

View file

@ -0,0 +1,29 @@
# Response caching
Design this crate for shared Rust execution used by the Python SDK and the Rust gateway. The Python SDK will remain, with more core execution moving to Rust and Python callbacks staying in Python. The Rust gateway is still evolving and is intended to replace the Python proxy. Keep response-cache policy independent of Python, HTTP serving, and either proxy's configuration format
Separate what is cached, how a hit is matched, and where entries are stored. Chat Completions, Messages, Responses, and embeddings are API workloads. Exact and semantic matching are lookup behaviors. Memory, Redis, disk, and object stores are storage choices. Embeddings are inference too, so do not use an inference-cache name to imply a category that excludes embeddings. Consult the existing Python cache and caching handler for behavior and compatibility contracts without copying their class structure
Storage traits, codecs, and backend capabilities belong in `litellm-cache` and the storage crates. Keep storage reusable for value types beyond LLM responses. This crate owns response entries, matching and freshness semantics, the Python-compatible response codec, and deferred-write policy. Core owns route-specific request identity, response encoding and reconstruction, embedding partial-hit orchestration, and stream capture and replay. Boundaries own configuration translation, resource construction, and caller identity
Construct and inject the response-cache service at the Python bridge or gateway boundary, as with the HTTP client. Reuse it across calls. Core and provider code must not discover cache configuration through Python globals, process configuration, or backend-specific factories
Keep `ResponseCache<B>` generic over its storage backend. Preserve typed backend contexts and capability bounds internally. Inject an object-safe service into core for runtime backend selection, so storage types do not spread through route and host types. Keep API request and response types statically typed. Add a generic parameter only where it preserves a useful type relationship or capability
Keep the core service contract narrow. Lookup and store must not require connection testing, ping, flush, deletion, counters, queues, or scripts. Require batch operations where a consumer needs partial hits, and keep management capabilities on their own interfaces. An exact-only adapter must remain explicit about its matching restriction. Supporting semantic matching requires a defined lookup-context and embedding execution contract, not just a renamed trait
Separate reusable resources from per-call policy. Backend configuration, namespace, default expiry, and entry limits belong to the configured service or backend. Read/write controls, expiry and freshness overrides, and authenticated caller scope belong to the call. Passing call options must not replace or mutate the route's configured service
Keep cache misses and storage failures distinguishable in return values. Core owns the decision to continue with provider execution after a cache failure. A read can reject an entry for freshness while the backend still retains it. Preserve the timestamp at which a response was produced when writing it later
Define lookup placement explicitly relative to authorization, deployment and credential resolution, and request-transforming callbacks. Cache identity must account for every input that affects reuse, including API surface and caller scope, while preserving intentional Python caching groups. Preserve existing keys and response formats unless changing them is an explicit migration decision
Cache normalized provider results before caller-specific response transformations. Hits must still run the applicable response processing, success callbacks, and cache-hit accounting. Keep callback execution in the host. Python cache implementations and semantic embedders that require the caller's task must use the existing host-operation mechanism rather than Python calls from a Rust worker. Preserve legacy fallback until that contract is supported
Keep unary caching independent of stream-only methods. Store streams only after successful exhaustion and protocol completion. Errors, incomplete streams, cancellation, and oversized entries must not populate the cache. Embedding batches need ordered partial results and reconstruction around the uncached inputs
Test each contract in its owner: storage capabilities in backend tests, envelopes and freshness here, reuse and replay in core, Python callback and fallback behavior at the bridge, and HTTP behavior at the gateway. Run backend contract checks and Python response-codec fixtures before exposing a new backend
`ScopedCache` requires an explicit shared or isolated scope at construction. Per-call `CachePolicy` controls reads, writes, expiry, and freshness without replacing the attached scope or service. `CacheOptions` binds that policy to an explicit scope for storage requests and has no default sharing policy. Versioned native envelopes reject incompatible API surfaces and versions as misses; this envelope is distinct from the legacy Python response codec
Response storage is not the source of budget or rate-limit coordination dependencies. Keep counters, reservations, and atomic admission operations out of `ResponseCacheService`, including when both services happen to use Redis

View file

@ -13,9 +13,12 @@ serde_json.workspace = true
sha2.workspace = true
[dev-dependencies]
litellm-cache-gcs.workspace = true
litellm-http = { workspace = true, features = ["test-support"] }
litellm-cache-memory.workspace = true
litellm-cache-redis.workspace = true
redis = "1.7.0"
redis-test = "1.0.4"
rstest.workspace = true
tokio.workspace = true
wiremock = "0.6.5"

View file

@ -1,51 +0,0 @@
# Response cache
`ResponseCache<B>` adds request keys, independent read/write controls, response envelopes, and freshness checks to any `B: BaseCache<Value = CacheEntry>`
## Ownership
`litellm-cache` defines typed storage, codec, and capability traits. `BaseCache` is only get, set, TTL, and pipeline writes. Everything else is an optional capability a backend implements only where its Python class defines the method: `DisconnectCache`, `ConnectionCache` (`test_connection`), `PingCache`, `BatchCache`, `DeleteCache`, `FlushCache`, counters, queues, TTL, scan, and scripts. Memory, Redis, disk, S3, GCS, and Azure Blob implement those traits without depending on response policy, so other consumers can store their own value types in the same backends
Semantic backends (Redis, Valkey, Qdrant) are generic over their embedder and codec, and share one prompt and embedding contract from `litellm_cache::semantic`. They take a `SemanticCacheContext`, so `ResponseCache` drives them the same way it drives exact backends
`litellm-cache-response` owns response keys, controls, entries, the Python-compatible response codec, and `WriteBuffer`, the backend-neutral deferred-write policy. It has no runtime dependency on a specific cache backend or Python
`ExactResponseCache` is the object-safe view of a `ResponseCache` over an exact backend. `ConnectionProbe` is the object-safe `test_connection`, implemented only when the backend implements `ConnectionCache`, so a host holds one next to its `ExactResponseCache` and reports the operation as unsupported otherwise, as Python's `BaseCache` does. Lookup, store, batch, and flush never require it
## Native Rust use
```rust
use std::{sync::Arc, time::Duration};
use litellm_cache_memory::InMemoryCache;
use litellm_cache_response::{CacheKeyInput, ResponseCache, ResponseCacheRequest};
use serde_json::json;
let cache = ResponseCache::new(Arc::new(InMemoryCache::default()));
let request = ResponseCacheRequest::new(CacheKeyInput {
preset: Some("example:key".into()),
..Default::default()
});
let now = Duration::from_secs(100);
cache.store(&request, json!({"answer": 7}), now)?;
assert_eq!(cache.async_lookup(&request, now).await?, Some(json!({"answer": 7})));
```
For Redis, inject `RedisCache::new(url, ttl, ResponseCacheCodec)` instead. Namespaces are optional and existing namespace prefixes are preserved
Callers supply Unix time for response freshness. Backend TTL uses its own clock. A read can reject an entry through `max_age` even while the backend still retains it
## Python integration
The bridge activates backends through the Rust catalog in `litellm/rust_bridge/catalog.py`. Every cache rule ships as `PYTHON_ONLY`, so SDK, Router, and proxy calls stay on Python and construct no native cache resources until a rule is changed
When a rule selects a backend, the Python `Cache` facade builds the native runtime from its own configuration and routes its storage calls (sync and async lookup and store, and pipelined batch store) to it. Stream replay, embedding partial-hit merging, response reconstruction, and callbacks stay in Python on top of that native store. The Python backend object remains for its direct API
Object responses are written as they are, and every other response shape is written as a serialized string, which is the pair of shapes Python reads. A string on the wire is therefore always a serialized response, so string-valued responses round trip. Typed backends such as memory never pass through the codec
Native cache handles must be recreated after fork. Native errors propagate to the host, which owns the existing fail-open and logging policy
## Adding another backend
Implement `BaseCache` for the backend with its associated value type and the capability traits its Python class supports, and accept a `CacheCodec` when wire serialization is needed. `ResponseCache<B>` then works without another response implementation
Run the `litellm-cache-testing` contract checks the backend's capabilities allow, and run response fixtures with `ResponseCacheCodec`, including both Python envelope encodings, before adding a catalog rule

View file

@ -4,6 +4,7 @@ mod codec;
mod embedding;
mod exact;
mod response;
mod service;
pub use buffer::WriteBuffer;
pub use caching::{
@ -14,3 +15,8 @@ pub use codec::ResponseCacheCodec;
pub use embedding::PartialHits;
pub use exact::{ConnectionProbe, ExactResponseCache};
pub use response::{ResponseCache, ResponseCacheRequest};
pub use service::{
CacheOptions, CachePolicy, CacheScope, ResponseCacheConfig, ResponseCacheService,
ResponseEnvelope, ScopedCache,
};

View file

@ -7,7 +7,9 @@ use litellm_cache::{
};
use serde_json::Value;
use crate::{CacheControls, CacheEntry, CacheKeyInput, PartialHits, cache_key};
use crate::{
CacheControls, CacheEntry, CacheKeyInput, PartialHits, ResponseCacheConfig, cache_key,
};
#[derive(Clone)]
pub struct ResponseCacheRequest<C: CacheContext = litellm_cache::ExactCacheContext> {
@ -50,6 +52,7 @@ where
B::Context: Default + PartialEq,
{
backend: Arc<B>,
config: ResponseCacheConfig,
}
impl<B> ResponseCache<B>
@ -58,7 +61,18 @@ where
B::Context: Default + PartialEq,
{
pub fn new(backend: Arc<B>) -> Self {
Self { backend }
Self {
backend,
config: ResponseCacheConfig::default(),
}
}
pub fn with_config(self, config: ResponseCacheConfig) -> Self {
Self { config, ..self }
}
pub fn config(&self) -> &ResponseCacheConfig {
&self.config
}
pub fn backend(&self) -> &B {
@ -221,7 +235,7 @@ where
response: Value,
now: Duration,
) -> Result<(), Error> {
if !request.controls.writes() {
if !request.controls.writes() || !self.fits(&response) {
return Ok(());
}
self.backend.set_cache(
@ -240,7 +254,7 @@ where
response: Value,
now: Duration,
) -> Result<(), Error> {
if !request.controls.writes() {
if !request.controls.writes() || !self.fits(&response) {
return Ok(());
}
self.backend
@ -277,7 +291,7 @@ where
) -> Result<(), Error> {
let writable = entries
.into_iter()
.filter(|(request, _, _)| request.controls.writes())
.filter(|(request, response, _)| request.controls.writes() && self.fits(response))
.map(|(request, response, now)| {
(
cache_key(&request.key),
@ -312,6 +326,11 @@ where
Ok(())
}
fn fits(&self, response: &Value) -> bool {
self.config.max_entry_bytes == usize::MAX
|| response.to_string().len() <= self.config.max_entry_bytes
}
fn partial_hits(
requests: &[ResponseCacheRequest<B::Context>],
readable: Vec<(usize, &ResponseCacheRequest<B::Context>)>,

View file

@ -0,0 +1,185 @@
use std::{future::Future, pin::Pin, time::Duration};
use litellm_cache::{BaseCache, Error, ExactCacheContext};
use serde_json::Value;
use crate::{
CacheControls, CacheEntry, CacheKeyField, CacheKeyInput, ResponseCache, ResponseCacheRequest,
};
type CacheFuture<'a, T> = Pin<Box<dyn Future<Output = Result<T, Error>> + Send + 'a>>;
#[derive(Clone)]
pub struct ResponseCacheConfig {
pub namespace: String,
pub max_entry_bytes: usize,
}
impl Default for ResponseCacheConfig {
fn default() -> Self {
Self {
namespace: String::new(),
max_entry_bytes: usize::MAX,
}
}
}
pub trait ResponseCacheService: Send + Sync {
fn config(&self) -> &ResponseCacheConfig;
fn lookup<'a>(
&'a self,
request: &'a ResponseCacheRequest,
now: Duration,
) -> CacheFuture<'a, Option<Value>>;
fn store<'a>(
&'a self,
request: &'a ResponseCacheRequest,
response: Value,
now: Duration,
) -> CacheFuture<'a, ()>;
}
impl<B> ResponseCacheService for ResponseCache<B>
where
B: BaseCache<Value = CacheEntry, Context = ExactCacheContext>,
{
fn config(&self) -> &ResponseCacheConfig {
self.config()
}
fn lookup<'a>(
&'a self,
request: &'a ResponseCacheRequest,
now: Duration,
) -> CacheFuture<'a, Option<Value>> {
Box::pin(self.async_lookup(request, now))
}
fn store<'a>(
&'a self,
request: &'a ResponseCacheRequest,
response: Value,
now: Duration,
) -> CacheFuture<'a, ()> {
Box::pin(self.async_store(request, response, now))
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum CacheScope {
Shared,
Isolated(String),
}
#[derive(Clone, Copy, Default)]
pub struct CachePolicy {
pub caching: Option<bool>,
pub no_cache: bool,
pub no_store: bool,
pub ttl: Option<Duration>,
pub max_age: Option<Duration>,
}
impl CachePolicy {
pub fn enabled(&self) -> bool {
self.caching != Some(false) && !(self.no_cache && self.no_store)
}
}
#[derive(Clone)]
pub struct CacheOptions {
pub policy: CachePolicy,
pub scope: CacheScope,
}
impl CacheOptions {
pub fn new(scope: CacheScope) -> Self {
Self {
policy: CachePolicy::default(),
scope,
}
}
pub fn request(self, namespace: &str, surface: &str, mut input: Value) -> ResponseCacheRequest {
input.sort_all_objects();
let scope = match self.scope {
CacheScope::Shared => String::new(),
CacheScope::Isolated(scope) => serde_json::json!(["isolated", scope]).to_string(),
};
ResponseCacheRequest {
key: CacheKeyInput {
namespace: Some(format!("{namespace}:inference-v2")),
fields: [
("surface", surface.to_owned()),
("scope", scope),
("request", input.to_string()),
]
.into_iter()
.map(|(name, value)| CacheKeyField {
name: name.into(),
value: Some(value),
api_parameter: true,
internal_parameter: false,
})
.collect(),
..Default::default()
},
controls: CacheControls {
configured: true,
supported_call_type: true,
native_backend: true,
default_on: true,
caching: self.policy.caching,
no_cache: self.policy.no_cache,
no_store: self.policy.no_store,
..Default::default()
},
context: ExactCacheContext {
ttl: self.policy.ttl,
},
max_age: self.policy.max_age,
}
}
}
#[derive(serde::Serialize, serde::Deserialize)]
pub struct ResponseEnvelope<T> {
version: u32,
surface: String,
output: T,
}
impl<T> ResponseEnvelope<T> {
pub fn new(surface: &str, output: T) -> Self {
Self {
version: 1,
surface: surface.into(),
output,
}
}
pub fn decode(self, surface: &str) -> Option<T> {
(self.version == 1 && self.surface == surface).then_some(self.output)
}
}
#[derive(Clone)]
pub struct ScopedCache {
pub service: std::sync::Arc<dyn ResponseCacheService>,
pub scope: CacheScope,
}
impl ScopedCache {
pub fn new(service: std::sync::Arc<dyn ResponseCacheService>, scope: CacheScope) -> Self {
Self { service, scope }
}
pub fn options(&self, policy: Option<CachePolicy>) -> CacheOptions {
CacheOptions {
policy: policy.unwrap_or_default(),
scope: self.scope.clone(),
}
}
}

View file

@ -18,7 +18,7 @@ use litellm_cache_response::{
WriteBuffer, cache_key,
};
use redis_test::MockCmd;
use rstest::rstest;
use rstest::{fixture, rstest};
use serde_json::{Value, json};
use support::{keyed, memory, redis, request};
@ -648,3 +648,129 @@ async fn write_buffer_clear_drops_pending_entries(memory: Memory, request: Respo
assert_eq!(memory.lookup(&request, now).unwrap(), None);
assert_eq!(memory.lookup(&other, now).unwrap(), None);
}
#[rstest]
#[case::python_sync("{'timestamp': 100.0, 'response': '{\"answer\": 7}'}")]
#[case::python_async(r#"{"timestamp":100.0,"response":{"answer":7}}"#)]
#[case::bare_response(r#"{"answer":7}"#)]
#[tokio::test]
async fn gcs_reads_python_entries_and_writes_python_compatible_envelopes(
#[case] encoded: &str,
#[values(false, true)] asynchronous: bool,
#[future(awt)] gcs: (wiremock::MockServer, Gcs),
) {
use wiremock::{
Mock, ResponseTemplate,
matchers::{body_json, header, method, path, query_param},
};
let (server, cache) = gcs;
let response = json!({"answer": 7});
Mock::given(method("GET"))
.and(path("/storage/v1/b/bucket/o/cache%2Fpython"))
.and(query_param("alt", "media"))
.and(header("authorization", "Bearer token"))
.respond_with(ResponseTemplate::new(200).set_body_string(encoded))
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/upload/storage/v1/b/bucket/o"))
.and(query_param("uploadType", "media"))
.and(query_param("name", "cache/native"))
.and(header("authorization", "Bearer token"))
.and(header("content-type", "application/json"))
.and(body_json(json!({"timestamp": 102.0, "response": response})))
.respond_with(ResponseTemplate::new(200))
.expect(1)
.mount(&server)
.await;
let lookup = if asynchronous {
cache
.async_lookup(&keyed("python"), Duration::from_secs(102))
.await
} else {
cache.lookup(&keyed("python"), Duration::from_secs(102))
};
assert_eq!(lookup.unwrap(), Some(response.clone()));
let request = ResponseCacheRequest {
context: litellm_cache::ExactCacheContext {
ttl: Some(Duration::from_secs(12)),
},
..keyed("native")
};
let stored = if asynchronous {
cache
.async_store(&request, response, Duration::from_secs(102))
.await
} else {
cache.store(&request, response, Duration::from_secs(102))
};
assert_eq!(stored, Ok(()));
let requests = server.received_requests().await.unwrap();
let upload = requests
.iter()
.find(|request| request.method.as_str() == "POST")
.unwrap();
assert_eq!(
upload.url.query(),
Some("uploadType=media&name=cache%2Fnative")
);
}
#[rstest]
#[tokio::test]
async fn gcs_batch_reads_preserve_order_and_treat_invalid_entries_as_misses(
#[future(awt)] gcs: (wiremock::MockServer, Gcs),
) {
use wiremock::{
Mock, ResponseTemplate,
matchers::{method, path},
};
let (server, cache) = gcs;
Mock::given(method("GET"))
.and(path("/storage/v1/b/bucket/o/cache%2Fhit"))
.respond_with(
ResponseTemplate::new(200)
.set_body_json(json!({"timestamp": 100.0, "response": {"answer":7}})),
)
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/storage/v1/b/bucket/o/cache%2Finvalid"))
.respond_with(ResponseTemplate::new(200).set_body_string("not an entry"))
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/storage/v1/b/bucket/o/cache%2Fmissing"))
.respond_with(ResponseTemplate::new(404))
.mount(&server)
.await;
let requests = [keyed("hit"), keyed("missing"), keyed("invalid")];
let partial = cache
.async_lookup_batch(&requests, Duration::from_secs(102))
.await
.unwrap();
assert_eq!(partial.values, vec![Some(json!({"answer":7})), None, None]);
assert_eq!(partial.missing_indices, vec![1, 2]);
}
type Gcs = ResponseCache<litellm_cache_gcs::GcsCache<litellm_cache_response::ResponseCacheCodec>>;
#[fixture]
async fn gcs() -> (wiremock::MockServer, Gcs) {
let server = wiremock::MockServer::start().await;
let cache = ResponseCache::new(Arc::new(litellm_cache_gcs::GcsCache::with_token_source(
litellm_cache_gcs::GcsConfig {
bucket_name: "bucket".into(),
gcs_path: Some("cache".into()),
path_service_account: None,
endpoint: server.uri(),
},
litellm_http::Client::plain_for_test(),
litellm_cache_response::ResponseCacheCodec,
Arc::new(litellm_cache_gcs::StaticTokenSource("token".into())),
)));
(server, cache)
}

View file

@ -0,0 +1,185 @@
use std::{
sync::{
Arc,
atomic::{AtomicU64, Ordering},
},
time::Duration,
};
use litellm_cache::ExactCacheContext;
use litellm_cache_memory::InMemoryCache;
use litellm_cache_response::{
CacheEntry, CacheKeyInput, ResponseCache, ResponseCacheConfig, ResponseCacheRequest,
ResponseCacheService,
};
use rstest::rstest;
use serde_json::json;
#[rstest]
#[tokio::test]
async fn service_honors_per_call_expiry_and_freshness() {
let clock = Arc::new(AtomicU64::new(0));
let cache_clock = clock.clone();
let cache: Arc<dyn ResponseCacheService> = Arc::new(ResponseCache::new(Arc::new(
InMemoryCache::with_clock(Some(100), Some(Duration::from_secs(60)), move || {
Duration::from_secs(cache_clock.load(Ordering::SeqCst))
}),
)));
let request = ResponseCacheRequest {
context: ExactCacheContext {
ttl: Some(Duration::from_secs(5)),
},
..ResponseCacheRequest::new(CacheKeyInput {
preset: Some("entry".into()),
..Default::default()
})
};
cache
.store(&request, json!({"answer":7}), Duration::ZERO)
.await
.unwrap();
assert_eq!(
cache.lookup(&request, Duration::ZERO).await.unwrap(),
Some(json!({"answer":7}))
);
let stale_request = ResponseCacheRequest {
max_age: Some(Duration::from_secs(1)),
..request.clone()
};
clock.store(2, Ordering::SeqCst);
assert_eq!(
cache
.lookup(&stale_request, Duration::from_secs(2))
.await
.unwrap(),
None
);
assert!(
cache
.lookup(&request, Duration::from_secs(2))
.await
.unwrap()
.is_some()
);
clock.store(6, Ordering::SeqCst);
assert_eq!(
cache
.lookup(&request, Duration::from_secs(6))
.await
.unwrap(),
None
);
}
#[rstest]
#[tokio::test]
async fn entry_limit_applies_to_sync_async_and_batch_writes() {
let storage = Arc::new(InMemoryCache::<CacheEntry>::default());
let cache = ResponseCache::new(storage.clone()).with_config(ResponseCacheConfig {
namespace: "service-test".into(),
max_entry_bytes: json!({"answer":7}).to_string().len(),
});
let small = json!({"answer":7});
let large = json!({"answer":"too large"});
let request = |key: &str| {
ResponseCacheRequest::new(CacheKeyInput {
preset: Some(key.into()),
..Default::default()
})
};
cache
.store(&request("sync"), large.clone(), Duration::ZERO)
.unwrap();
cache
.async_store(&request("async"), large.clone(), Duration::ZERO)
.await
.unwrap();
cache
.async_store_batch(
vec![
(request("batch-large"), large),
(request("batch-small"), small.clone()),
],
Duration::ZERO,
)
.await
.unwrap();
let service: Arc<dyn ResponseCacheService> = Arc::new(cache);
service
.store(&request("service"), small.clone(), Duration::ZERO)
.await
.unwrap();
for key in ["sync", "async", "batch-large"] {
assert!(storage.get_cache(key).unwrap().is_none());
}
for key in ["batch-small", "service"] {
assert_eq!(
service.lookup(&request(key), Duration::ZERO).await.unwrap(),
Some(small.clone())
);
}
}
#[rstest]
#[case::same_scope("tenant-a", "tenant-a", true)]
#[case::different_scope("tenant-a", "tenant-b", false)]
#[case::empty_isolated_scope("", "", true)]
#[tokio::test]
async fn isolated_policy_controls_actual_entry_reuse(
#[case] first: &str,
#[case] second: &str,
#[case] hit: bool,
#[values(false, true)] override_policy: bool,
) {
use litellm_cache_response::{CachePolicy, CacheScope, ScopedCache};
let service = Arc::new(ResponseCache::new(Arc::new(
InMemoryCache::<CacheEntry>::default(),
)));
let request = |scope| {
ScopedCache::new(service.clone(), scope)
.options(override_policy.then_some(CachePolicy {
ttl: Some(Duration::from_secs(30)),
..CachePolicy::default()
}))
.request("test", "messages", json!({"prompt":"hello"}))
};
service
.async_store(
&request(CacheScope::Isolated(first.into())),
json!({"answer":7}),
Duration::ZERO,
)
.await
.unwrap();
assert_eq!(
service
.async_lookup(
&request(CacheScope::Isolated(second.into())),
Duration::ZERO
)
.await
.unwrap(),
hit.then(|| json!({"answer":7}))
);
assert_eq!(
service
.async_lookup(&request(CacheScope::Shared), Duration::ZERO)
.await
.unwrap(),
None
);
}
#[rstest]
#[case::valid(1, "messages", Some(7))]
#[case::unknown_version(2, "messages", None)]
#[case::another_surface(1, "responses", None)]
fn envelopes_require_a_matching_surface_and_version(
#[case] version: u32,
#[case] surface: &str,
#[case] expected: Option<u32>,
) {
let envelope: litellm_cache_response::ResponseEnvelope<u32> =
serde_json::from_value(json!({"version":version,"surface":surface,"output":7})).unwrap();
assert_eq!(envelope.decode("messages"), expected);
}

View file

@ -1,12 +1,13 @@
- Target invariants, not completion claims
- This crate is the legacy `@client` wrapper as the native call sees it, and nothing else: the `Logging` contract (`function_setup`, the deployment hooks, `pre_call`/`post_call`, the sync and async success and failure fan-out, the deferred proxy release, the argument sharing those callbacks rely on)
- This crate owns compatibility for all existing Python callbacks and loggers, including `CustomLogger`. `mapping.rs` owns the executable call bindings and the inventory of Python-owned hooks. A Python-owned entry records an existing path, never permission to invoke it a second time. The native call adapter preserves the `Logging` contract (`function_setup`, the deployment hooks, `pre_call`/`post_call`, the sync and async success and failure fan-out, the deferred proxy release, the argument sharing those callbacks rely on)
- Smell test: if a future callback host (`callbacks-v1-python`, WASM, in-process Rust) could share a piece of this crate, it does not belong here
- SDK request policy (credential inheritance, the budget and retry-count limits) is the driver's preflight, supplied by `python-bridge`; this crate only adopts the keyword view it produces
- The driver in `litellm-host-python`, the routes and core see one `PythonLifecycle`; they never learn which Python objects consume a call
- SDK request policy (credential inheritance, the budget and retry-count limits) is a separate hook supplied by `python-bridge`; compose it after this adapter so logging adopts the final keyword view before policy mutates or rejects it
- The driver in `litellm-host-python`, the routes and core see one `PythonCallHooks` using the shared `CallEvent`; they never learn which Python objects consume a call
- Every litellm Python internal Rust still borrows is a variant of `LegacyPython`, grouped by subsystem, with its signature pinned in `python_contract.json`
- The enum only shrinks: when Rust owns a subsystem, delete its group rather than adding a Rust path beside it
- Calling a user's own callback directly is permanent Python surface and gets its own type outside `LegacyPython`
- `PublicCall` is the caller's call as `Logging` sees it: the positional arguments, the keyword view as the call rewrites it (setup, deployment hook, preflight) and the bound request object backing omitted keywords; routes hand it over through `run_legacy_call` and keep no copy
- `PublicCall` is the caller's call as `Logging` sees it: the positional arguments, the keyword view as the call rewrites it (setup, deployment hook, preflight) and the bound request object backing omitted keywords; shared bridge composition hands it to `LegacyLogging`; routes use the neutral call boundary
- `LoggingOperation` selects legacy logging entrypoints and response handling. It belongs here rather than in shared inference data contracts
- `setup` reuses a `Logging` passed as `litellm_logging_obj` (the proxy and Router) and otherwise builds one through `function_setup`; which callbacks run is `Logging`'s decision, never this crate's
- Callbacks receive the caller's own objects and may mutate them; this crate alone carries that obligation
- Retain complete boundary arguments, opaque values, aliases, omitted/default distinctions and deliberate copies; preserve the deployment-hook kwargs view

File diff suppressed because it is too large Load diff

View file

@ -3,16 +3,13 @@
//! lifetime. No other callback host has that obligation, which is why nothing outside
//! this crate holds them.
use litellm_host::{machine::Machine, protocol::Protocol};
use litellm_host_python::{Preflight, ProtocolHost, lookup, run_call};
use litellm_host_python::lookup;
use pyo3::{
gc::{PyTraverseError, PyVisit},
prelude::*,
types::{PyDict, PyTuple},
};
use crate::{LegacyLogging, LegacySurface};
pub struct PublicCall {
args: Py<PyTuple>,
kwargs: Py<PyDict>,
@ -34,6 +31,10 @@ impl PublicCall {
})
}
pub fn arguments(&self, py: Python<'_>) -> Py<PyDict> {
self.kwargs.clone_ref(py)
}
pub(crate) fn args(&self) -> &Py<PyTuple> {
&self.args
}
@ -64,34 +65,6 @@ impl PublicCall {
}
}
/// Runs one native call under the legacy `Logging` contract: the protocol host projects from
/// the keyword view the contract prepares and `preflight` rewrites, and the contract
/// observes the call.
pub fn run_legacy_call<H, M>(
py: Python<'_>,
surface: LegacySurface,
call: PublicCall,
machine: M,
host: H,
preflight: Preflight,
asynchronous: bool,
) -> PyResult<Py<PyAny>>
where
H: ProtocolHost + 'static,
M: Machine<Protocol = H::Protocol, Complete = <H::Protocol as Protocol>::Response> + 'static,
{
let arguments = call.kwargs.clone_ref(py);
run_call(
py,
machine,
host,
Box::new(LegacyLogging::new(py, surface, call, asynchronous)),
preflight,
arguments,
asynchronous,
)
}
#[cfg(test)]
mod tests {
use super::*;
@ -110,7 +83,7 @@ mod tests {
(call, locals)
}
#[test]
#[rstest::rstest]
fn capture_copies_the_keyword_dict_without_copying_its_values() {
Python::initialize();
Python::attach(|py| {

View file

@ -2,7 +2,7 @@
//! the deferred and worker-submitted success paths, and the sync-callbacks-for-async-calls
//! duplication. All of it expires with the legacy callback contract.
use litellm_host::event::{RequestContext, WireRequest};
use litellm_host::interceptors::{RequestContext, WireRequest};
use litellm_host_python::to_py;
use pyo3::{exceptions::PyBaseException, prelude::*, types::PyDict};

View file

@ -1,26 +1,32 @@
//! The legacy `@client` wrapper as the native call sees it: litellm's `Logging` object, the
//! sync and async callback registries it fans out to, the deployment hooks and the deferred
//! proxy release. All of it sits behind one
//! [`PythonLifecycle`](litellm_host_python::PythonLifecycle), so the driver, the routes and
//! core never learn which Python object is on the other end. The SDK's own request policy
//! (credential inheritance, the budget and retry limits) is the driver's preflight, not this
//! crate's.
//! [`PythonCallHooks`](litellm_host_python::PythonCallHooks), so the driver, the routes and
//! core never learn which Python object is on the other end.
//!
//! Legacy callbacks receive the caller's own objects and may mutate them. [`PublicCall`]
//! is where those objects live, and [`run_legacy_call`] is how a route hands them over
//! without keeping a copy.
//! is where those objects live.
mod adapter;
mod call;
mod callbacks;
mod deferred;
mod logger;
mod mapping;
mod python;
pub(crate) use adapter::LegacyLogging;
pub use adapter::{LegacySurface, PassThroughStream};
pub use call::{PublicCall, run_legacy_call};
pub use adapter::LegacyLogging;
pub use call::PublicCall;
pub(crate) use callbacks::{LegacyCallbacks, is_internal_call};
pub(crate) use logger::{DeploymentHooks, PythonLogger, finalize, setup};
pub use mapping::{CallBoundary, CallbackMapping, Dispatch, callback_mappings};
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum LoggingOperation {
Completion,
Responses,
Messages,
Ocr,
}
#[cfg(test)]
mod test_support;

View file

@ -0,0 +1,288 @@
use litellm_host::{
hooks::CallHooks,
interceptors::{RawResponse, RequestContext, WireRequest},
lifecycle::{ExecutionEvent, FailureOrigin, Timing},
};
use litellm_host_python::{HookStep, PythonCallEvent, PythonRuntime};
use pyo3::{prelude::*, types::PyDict};
use crate::LegacyLogging;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum CallBoundary {
PrepareArguments,
BeforeProviderRequest,
AfterProviderResponse,
TransformResponse,
Succeeded,
Failed,
StreamOpened,
StreamChunk,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Dispatch {
Call(CallBoundary),
Python(&'static str),
DeclarationOnly,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct CallbackMapping {
pub callback: &'static str,
pub dispatch: Dispatch,
}
struct Binding<H> {
boundary: CallBoundary,
invoke: H,
callbacks: &'static [&'static str],
}
impl<H> Binding<H> {
fn mappings(&self) -> impl Iterator<Item = CallbackMapping> {
self.callbacks.iter().map(|callback| CallbackMapping {
callback,
dispatch: Dispatch::Call(self.boundary),
})
}
}
type Step<T> = PyResult<HookStep<LegacyLogging, T>>;
type Prepare = fn(&mut LegacyLogging, Python<'_>, Py<PyDict>, f64) -> Step<Py<PyDict>>;
type Before =
fn(&mut LegacyLogging, Python<'_>, Box<WireRequest>, &RequestContext) -> Step<Box<WireRequest>>;
type After = fn(&mut LegacyLogging, Python<'_>, &RawResponse) -> Step<()>;
type Transform = fn(&mut LegacyLogging, Python<'_>, Py<PyAny>, Timing) -> Step<Py<PyAny>>;
type Success = fn(&mut LegacyLogging, Python<'_>, Timing, &Py<PyAny>) -> Step<()>;
type Failure = fn(&mut LegacyLogging, Python<'_>, Timing, FailureOrigin, &PyErr) -> Step<()>;
type Open = fn(&mut LegacyLogging, Python<'_>, &Py<PyAny>) -> PyResult<()>;
type Chunk = fn(&mut LegacyLogging, Python<'_>, &Py<PyAny>) -> PyResult<()>;
const PREPARE: Binding<Prepare> = Binding {
boundary: CallBoundary::PrepareArguments,
invoke: LegacyLogging::prepare_call,
callbacks: &["async_pre_call_deployment_hook"],
};
const BEFORE: Binding<Before> = Binding {
boundary: CallBoundary::BeforeProviderRequest,
invoke: LegacyLogging::pre_call,
callbacks: &["log_pre_api_call", "log_input_event"],
};
const AFTER: Binding<After> = Binding {
boundary: CallBoundary::AfterProviderResponse,
invoke: LegacyLogging::post_call,
callbacks: &["log_post_api_call"],
};
const TRANSFORM: Binding<Transform> = Binding {
boundary: CallBoundary::TransformResponse,
invoke: LegacyLogging::transform_public_response,
callbacks: &["async_post_call_success_deployment_hook"],
};
const SUCCESS: Binding<Success> = Binding {
boundary: CallBoundary::Succeeded,
invoke: LegacyLogging::succeeded,
callbacks: &[
"log_success_event",
"async_log_success_event",
"logging_hook",
"async_logging_hook",
"redact_standard_logging_payload_from_model_call_details",
"log_event",
"async_log_event",
],
};
const FAILURE: Binding<Failure> = Binding {
boundary: CallBoundary::Failed,
invoke: LegacyLogging::failed,
callbacks: &[
"async_post_call_failure_deployment_hook",
"log_failure_event",
"async_log_failure_event",
"log_model_group_rate_limit_error",
"log_event",
"async_log_event",
],
};
const OPEN: Binding<Open> = Binding {
boundary: CallBoundary::StreamOpened,
invoke: LegacyLogging::stream_opened,
callbacks: &[],
};
const CHUNK: Binding<Chunk> = Binding {
boundary: CallBoundary::StreamChunk,
invoke: LegacyLogging::stream_chunk,
callbacks: &[],
};
pub fn callback_mappings() -> impl Iterator<Item = CallbackMapping> {
PREPARE
.mappings()
.chain(BEFORE.mappings())
.chain(AFTER.mappings())
.chain(TRANSFORM.mappings())
.chain(SUCCESS.mappings())
.chain(FAILURE.mappings())
.chain(OPEN.mappings())
.chain(CHUNK.mappings())
.chain(PYTHON_CALLBACKS.iter().copied())
}
impl CallHooks<PythonRuntime> for LegacyLogging {
fn prepare_arguments(
&mut self,
py: Python<'_>,
arguments: Py<PyDict>,
started_at: f64,
) -> Step<Py<PyDict>> {
(PREPARE.invoke)(self, py, arguments, started_at)
}
fn arguments_prepared(&mut self, py: Python<'_>, arguments: &Py<PyDict>) -> PyResult<()> {
self.adopt_arguments(py, arguments);
Ok(())
}
fn before_provider_request(
&mut self,
py: Python<'_>,
wire: Box<WireRequest>,
context: &RequestContext,
) -> Step<Box<WireRequest>> {
(BEFORE.invoke)(self, py, wire, context)
}
fn transform_response(
&mut self,
py: Python<'_>,
response: Py<PyAny>,
timing: Timing,
) -> Step<Py<PyAny>> {
(TRANSFORM.invoke)(self, py, response, timing)
}
fn on_event(&mut self, py: Python<'_>, event: PythonCallEvent<'_>) -> Step<()> {
match event {
PythonCallEvent::Started { .. } | PythonCallEvent::Cancelled { .. } => {
Ok(HookStep::Ready(()))
}
PythonCallEvent::Execution(ExecutionEvent::ResultReady { facts }) => {
self.result_ready(py, &facts)
}
PythonCallEvent::Execution(ExecutionEvent::ProviderResponseReceived { raw }) => {
(AFTER.invoke)(self, py, raw)
}
PythonCallEvent::Succeeded { timing, response } => {
(SUCCESS.invoke)(self, py, timing, response)
}
PythonCallEvent::Failed {
timing,
origin,
error,
} => (FAILURE.invoke)(self, py, timing, origin, error),
}
}
fn on_stream_open(&mut self, py: Python<'_>, head: &Py<PyAny>) -> PyResult<()> {
(OPEN.invoke)(self, py, head)
}
fn on_stream_chunk(&mut self, py: Python<'_>, chunk: &Py<PyAny>) -> PyResult<()> {
(CHUNK.invoke)(self, py, chunk)
}
}
macro_rules! python_callbacks {
($($dispatch:expr => [$($callback:literal),* $(,)?]),* $(,)?) => {
const PYTHON_CALLBACKS: &[CallbackMapping] = &[
$($(CallbackMapping { callback: $callback, dispatch: $dispatch },)*)*
];
};
}
python_callbacks! {
Dispatch::Python("litellm.router") => [
"async_pre_routing_hook",
"async_filter_deployments",
"pre_call_check",
"async_pre_call_check",
],
Dispatch::Python("litellm.router_utils.fallback_event_handlers") => [
"log_success_fallback_event",
"log_failure_fallback_event",
],
Dispatch::Python("litellm.proxy.utils") => [
"async_pre_call_hook",
"async_post_call_response_headers_hook",
"async_post_call_failure_hook",
"async_post_call_success_hook",
"async_moderation_hook",
"async_post_call_streaming_hook",
"async_post_call_streaming_iterator_hook",
"async_filter_listed_models",
],
Dispatch::Python("litellm.litellm_core_utils.litellm_logging") => [
"async_get_chat_completion_prompt",
"get_chat_completion_prompt",
"log_stream_event",
"async_log_stream_event",
"async_post_mcp_tool_call_hook",
],
Dispatch::Python("litellm.llms.anthropic.pass_through.messages.handler") => [
"async_pre_request_hook",
],
Dispatch::Python("litellm.litellm_core_utils.streaming_handler") => [
"async_post_call_streaming_deployment_hook",
],
Dispatch::Python("litellm.responses.streaming_iterator") => [
"async_post_call_streaming_deployment_hook",
],
Dispatch::Python("litellm.main") => [
"translate_completion_input_params",
"translate_completion_output_params",
"translate_completion_output_params_streaming",
],
Dispatch::Python("litellm.integrations.argilla") => ["async_dataset_hook"],
Dispatch::Python("litellm.proxy.management_helpers.audit_logs") => ["async_log_audit_log_event"],
Dispatch::Python("litellm.llms.custom_httpx.llm_http_handler") => [
"async_should_run_agentic_loop",
"async_run_agentic_loop",
"async_build_agentic_loop_plan",
"async_post_agentic_loop_response_hook",
"async_agentic_loop_cleanup_hook",
"async_should_run_chat_completion_agentic_loop",
"async_run_chat_completion_agentic_loop",
"async_build_chat_completion_agentic_loop_plan",
],
Dispatch::Python("litellm.litellm_core_utils.chat_completion_agentic_loop") => [
"async_should_run_agentic_loop",
"async_run_agentic_loop",
"async_build_agentic_loop_plan",
"async_post_agentic_loop_response_hook",
"async_agentic_loop_cleanup_hook",
],
Dispatch::Python("litellm.llms.openai.openai") => [
"async_should_run_chat_completion_agentic_loop",
"async_run_chat_completion_agentic_loop",
],
Dispatch::Python("litellm.proxy.spend_tracking.cold_storage_handler") => [
"get_proxy_server_request_from_cold_storage_with_object_key",
],
Dispatch::Python("litellm.integrations.custom_logger") => [
"truncate_standard_logging_payload_content",
"redacts_messages_itself",
"handle_callback_failure",
"get_callback_env_vars",
],
Dispatch::DeclarationOnly => [
"async_log_pre_api_call",
"async_log_input_event",
],
}

Some files were not shown because too many files have changed in this diff Show more