diff --git a/.circleci/scripts/run_integration.sh b/.circleci/scripts/run_integration.sh
index ba24e66ba1c..f3166e5db39 100644
--- a/.circleci/scripts/run_integration.sh
+++ b/.circleci/scripts/run_integration.sh
@@ -173,7 +173,7 @@ start_proxy() {
AWS_EC2_METADATA_DISABLED=true DO_NOT_TRACK=1 COVERAGE_FILE="$coverage_data" \
"${proxy_command[@]}" --config tests/integration/proxy_config.yaml \
--host 127.0.0.1 --port "$port" --num_workers 1 --telemetry False \
- --use_prisma_db_push --enforce_prisma_migration_check \
+ --use_prisma_db_push \
> "$results/$log_name" 2>&1 &
launched_pid=$!
}
diff --git a/.circleci/scripts/unit_selection.sh b/.circleci/scripts/unit_selection.sh
index bfaa27c3ef3..542984dd2e0 100755
--- a/.circleci/scripts/unit_selection.sh
+++ b/.circleci/scripts/unit_selection.sh
@@ -77,6 +77,7 @@ legacy_paths() {
echo tests/unit/embeddings
echo tests/unit/endpoints
echo tests/unit/files
+ echo tests/unit/harness
echo tests/unit/images
echo tests/unit/interactions
echo tests/unit/messages
@@ -145,6 +146,7 @@ legacy_paths() {
echo tests/unit/proxy/test_proxy_token_counter.py
echo tests/unit/proxy/test_server_root_path.py ;;
proxy-db-proxy-server-core)
+ echo tests/unit/proxy/test__lazy_features.py
echo tests/unit/proxy/test_aproxy_startup.py
echo tests/unit/proxy/test_proxy_server.py ;;
proxy-db-proxy-utils) echo tests/unit/proxy/test_proxy_utils.py ;;
diff --git a/.github/ci-coverage-allowlist.yml b/.github/ci-coverage-allowlist.yml
index 445a8519436..eea25e8e285 100644
--- a/.github/ci-coverage-allowlist.yml
+++ b/.github/ci-coverage-allowlist.yml
@@ -4,6 +4,14 @@ description: >-
by a job nor listed here, so every entry below is a decision on the record.
test_paths:
+ - reason: >-
+ litellm.agent() end-to-end suite. It drives the real claude, codex and opencode CLIs and
+ deepagents against a live LiteLLM AI Gateway, so it needs those binaries on PATH plus
+ LITELLM_PROXY_API_BASE / LITELLM_PROXY_API_KEY, and skips without them. Run manually
+ before changing litellm/harness; the mocked coverage runs in tests/unit/harness and
+ tests/unit/llms/*/harness
+ paths:
+ - tests/harness_e2e
- reason: >-
The Rust/Python parity harness is run manually through its local CLI. Recorded replay,
fixture generation, and harness checks are intentionally outside pull request CI
diff --git a/README.md b/README.md
index 98c5343daee..4004e6474ee 100644
--- a/README.md
+++ b/README.md
@@ -268,6 +268,31 @@ For MCP OAuth, an upstream may advertise dynamic client registration but refuse
+
+Agents - Run Claude Code, Codex, OpenCode or Deep Agents on any model (Python SDK)
+
+### Python SDK - Agents
+
+```python
+import litellm
+from litellm import Harness, sandbox
+
+result = litellm.agent(
+ Harness.CLAUDE_CODE, # or Harness.CODEX, Harness.OPENCODE, Harness.DEEPAGENTS
+ "Find why tests/test_router.py is flaky and fix it.",
+ sandbox=sandbox.local("./repo"),
+ model="litellm_proxy/claude-sonnet-4-5", # a model group on your AI Gateway
+)
+
+print(result.text, result.cost, [f.path for f in result.files])
+```
+
+Set `LITELLM_PROXY_API_BASE` and `LITELLM_PROXY_API_KEY` and every model call the agent makes goes through your AI Gateway, tagged `harness,claude_code`. Drop the `litellm_proxy/` prefix to call a provider directly. Install `starlette uvicorn` plus the agent's CLI (`claude`, `codex` or `opencode`), or `deepagents langchain-litellm` for Deep Agents.
+
+[**Docs: Agent Harnesses**](https://docs.litellm.ai/docs/harness)
+
+
+
### Supported Providers ([Website Supported Models](https://models.litellm.ai/) | [Docs](https://docs.litellm.ai/docs/providers))
| Provider | `/chat/completions` | `/messages` | `/responses` | `/embeddings` | `/image/generations` | `/audio/transcriptions` | `/audio/speech` | `/moderations` | `/batches` | `/rerank` |
diff --git a/backend/routes/allowlist.py b/backend/routes/allowlist.py
index 80ca0ef22bb..d7f3e615c67 100644
--- a/backend/routes/allowlist.py
+++ b/backend/routes/allowlist.py
@@ -60,6 +60,7 @@ BACKEND_PATH_PREFIXES: tuple[str, ...] = (
# Tools / agents (registry & policy admin)
"/v1/tool/",
"/v1/agents",
+ "/agent/daily/activity/",
# Guardrails admin
"/v2/guardrails/",
# MCP server admin + BYOK OAuth flow (UI-initiated) + dynamic per-server endpoints
diff --git a/deploy/lens/compose.yaml b/deploy/lens/compose.yaml
index d41cb8eb203..4d1224fd41e 100644
--- a/deploy/lens/compose.yaml
+++ b/deploy/lens/compose.yaml
@@ -1,6 +1,6 @@
services:
lens-worker:
- image: ${LENS_WORKER_IMAGE:-ghcr.io/berriai/litellm-lens-worker@sha256:a8e8731d954916594eea462969946b9292fb771681ff515a9fd296b53f856c77}
+ image: ${LENS_WORKER_IMAGE:-ghcr.io/berriai/litellm-lens-worker@sha256:67eba741c1b97c749975c5c38e2370a603e1105babc908d613c1b79d7b995393}
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}
diff --git a/docker-compose.hardened.yml b/docker-compose.hardened.yml
index 31d0c2e9ef2..84a23faa054 100644
--- a/docker-compose.hardened.yml
+++ b/docker-compose.hardened.yml
@@ -6,8 +6,6 @@ services:
context: .
dockerfile: docker/Dockerfile.non_root
target: runtime
- args:
- PROXY_EXTRAS_SOURCE: "local"
depends_on:
- squid
user: "101:101"
diff --git a/docker/Dockerfile.non_root b/docker/Dockerfile.non_root
index eca12855afa..ca526e06834 100644
--- a/docker/Dockerfile.non_root
+++ b/docker/Dockerfile.non_root
@@ -3,7 +3,6 @@
# Base images
ARG LITELLM_BUILD_IMAGE=cgr.dev/chainguard/wolfi-base@sha256:1d95114038f76513a9ace6fca107d5582b08c65981f81f61cb56bf7fd2ef216d
ARG LITELLM_RUNTIME_IMAGE=cgr.dev/chainguard/wolfi-base@sha256:1d95114038f76513a9ace6fca107d5582b08c65981f81f61cb56bf7fd2ef216d
-ARG PROXY_EXTRAS_SOURCE=published
ARG UV_IMAGE=ghcr.io/astral-sh/uv:0.11.7@sha256:240fb85ab0f263ef12f492d8476aa3a2e4e1e333f7d67fbdd923d00a506a516a
# Pinned by digest like the other base images; bump explicitly on Node upgrades.
ARG UI_BUILD_IMAGE=node:24.19-alpine3.24@sha256:d32cdf619f63fe0471182d08996dd516c6275bb5fd31ae06e55a570bd9e1ad43
@@ -44,7 +43,6 @@ COPY ui/litellm-dashboard/ ./
RUN npm run build
FROM $LITELLM_BUILD_IMAGE AS builder
-ARG PROXY_EXTRAS_SOURCE
WORKDIR /app
USER root
@@ -107,26 +105,14 @@ RUN mkdir -p /var/lib/litellm/ui /var/lib/litellm/assets && \
touch /var/lib/litellm/ui/.litellm_ui_ready
RUN --mount=type=cache,target=/app/.cache/uv,id=litellm-uv-cache \
- if [ "$PROXY_EXTRAS_SOURCE" = "published" ]; then \
- uv sync --frozen --no-default-groups --no-editable \
- --extra proxy \
- --extra proxy-runtime \
- --extra extra_proxy \
- --extra semantic-router \
- --extra saml \
- --extra bedrock-realtime \
- --python python3.13 \
- --no-sources-package litellm-proxy-extras; \
- else \
- uv sync --frozen --no-default-groups --no-editable \
- --extra proxy \
- --extra proxy-runtime \
- --extra extra_proxy \
- --extra semantic-router \
- --extra saml \
- --extra bedrock-realtime \
- --python python3.13; \
- fi
+ uv sync --frozen --no-default-groups --no-editable \
+ --extra proxy \
+ --extra proxy-runtime \
+ --extra extra_proxy \
+ --extra semantic-router \
+ --extra saml \
+ --extra bedrock-realtime \
+ --python python3.13
RUN HOME=/opt/prisma XDG_CACHE_HOME=/opt/prisma/.cache PRISMA_BINARY_CACHE_DIR=/opt/prisma/binaries \
npm_config_cache=/root/.npm \
@@ -136,7 +122,6 @@ RUN sed -i 's/\r$//' docker/entrypoint.sh && chmod +x docker/entrypoint.sh && \
sed -i 's/\r$//' docker/prod_entrypoint.sh && chmod +x docker/prod_entrypoint.sh
FROM $LITELLM_RUNTIME_IMAGE AS runtime
-ARG PROXY_EXTRAS_SOURCE
WORKDIR /app
USER root
diff --git a/enterprise/pyproject.toml b/enterprise/pyproject.toml
index 74cedb9d84d..43aa5a1f728 100644
--- a/enterprise/pyproject.toml
+++ b/enterprise/pyproject.toml
@@ -1,6 +1,6 @@
[project]
name = "litellm-enterprise"
-version = "0.1.72"
+version = "0.1.73"
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.72"
+version = "0.1.73"
version_files = [
"pyproject.toml:^version",
"../pyproject.toml:litellm-enterprise==",
diff --git a/litellm-proxy-extras/litellm_proxy_extras/migration_lock.py b/litellm-proxy-extras/litellm_proxy_extras/migration_lock.py
index e4ccbe585a9..bea5e36fd18 100644
--- a/litellm-proxy-extras/litellm_proxy_extras/migration_lock.py
+++ b/litellm-proxy-extras/litellm_proxy_extras/migration_lock.py
@@ -87,3 +87,21 @@ def migration_lock(database_url: str) -> Generator[MigrationCoordinator, None, N
f"Timed out waiting for another v2 migration resolver after {wait_seconds}s. "
f"Check the running migration or increase {MIGRATION_LOCK_TIMEOUT_ENV_VAR}."
)
+
+
+@contextmanager
+def held_migration_lock(connection: "psycopg.Connection[tuple[object, ...]]") -> Generator[bool, None, None]:
+ """A session-level, non-blocking hold of the migration coordinator lock on an autocommit
+ connection, for DDL that cannot run inside a transaction (`CREATE INDEX CONCURRENTLY`).
+ Yields whether the lock was acquired; a v2 resolver or another migration job's index build
+ holding it yields False. Released on exit."""
+ from psycopg.rows import class_row
+
+ with connection.cursor(row_factory=class_row(_LockResult)) as cursor:
+ row: Final = cursor.execute("SELECT pg_try_advisory_lock(%s) AS acquired", (MIGRATION_LOCK_KEY,)).fetchone()
+ acquired: Final = row is not None and row.acquired
+ try:
+ yield acquired
+ finally:
+ if acquired:
+ connection.execute("SELECT pg_advisory_unlock(%s)", (MIGRATION_LOCK_KEY,))
diff --git a/litellm-proxy-extras/litellm_proxy_extras/migration_recovery.py b/litellm-proxy-extras/litellm_proxy_extras/migration_recovery.py
index 9202317c776..5a55b35b255 100644
--- a/litellm-proxy-extras/litellm_proxy_extras/migration_recovery.py
+++ b/litellm-proxy-extras/litellm_proxy_extras/migration_recovery.py
@@ -1,4 +1,5 @@
import hashlib
+import re
import subprocess
from collections.abc import Mapping
from dataclasses import dataclass
@@ -156,3 +157,48 @@ def baseline_current_schema(
"review any feature-specific backfill requirements.",
len(migrations),
)
+
+
+_LINE_COMMENT_RE: Final = re.compile(r"--[^\n]*")
+_BLOCK_COMMENT_RE: Final = re.compile(r"/\*.*?\*/", re.DOTALL)
+_NO_OP_STATEMENT_RE: Final = re.compile(r"^\s*SELECT\s+1\s*$", re.IGNORECASE)
+
+
+def is_inert_migration(script: str) -> bool:
+ """Whether a migration file changes nothing: only comments and `SELECT 1`, so
+ applying it can neither repeat nor skip a database change."""
+ stripped: Final = _LINE_COMMENT_RE.sub("", _BLOCK_COMMENT_RE.sub("", script))
+ return all(not part.strip() or _NO_OP_STATEMENT_RE.match(part) for part in stripped.split(";"))
+
+
+def roll_back_failed_inert_migration(coordinator: MigrationCoordinator, schema: str, migration: Path) -> bool:
+ """Roll back the failed ledger row of a migration whose file in this build is inert,
+ so `migrate deploy` applies the inert file on its next pass. The row records an
+ earlier build's attempt at SQL this build no longer ships (an index now built by the
+ migration job), so no database change can be repeated or skipped by replaying
+ the empty file. The caller commits this checkpoint before the next Prisma command.
+ """
+ from psycopg import sql
+
+ if not is_inert_migration(migration.read_text(encoding="utf-8")):
+ return False
+ coordinator.acquire_prisma_lock()
+ records: Final = _migration_records(coordinator.connection, schema, migration)
+ unfinished: Final = tuple(record for record in records if not record.finished)
+ if len(unfinished) != 1:
+ return False
+ result: Final = coordinator.connection.execute(
+ sql.SQL(
+ "UPDATE {} SET rolled_back_at = current_timestamp "
+ "WHERE id = %s AND finished_at IS NULL AND rolled_back_at IS NULL"
+ ).format(sql.Identifier(schema, "_prisma_migrations")),
+ (unfinished[0].id,),
+ )
+ if result.rowcount != 1:
+ raise RuntimeError("Could not roll back the failed inert migration history row; rerun the database setup.")
+ logger.info(
+ "Rolled back the failed history row of %s: this build ships it as an inert migration, "
+ "its index is built by the migration job",
+ migration.parent.name,
+ )
+ return True
diff --git a/litellm-proxy-extras/litellm_proxy_extras/migrations/20260823000000_add_spend_logs_api_key_starttime_index/migration.sql b/litellm-proxy-extras/litellm_proxy_extras/migrations/20260823000000_add_spend_logs_api_key_starttime_index/migration.sql
index 9a061aaed43..a2bec81ca00 100644
--- a/litellm-proxy-extras/litellm_proxy_extras/migrations/20260823000000_add_spend_logs_api_key_starttime_index/migration.sql
+++ b/litellm-proxy-extras/litellm_proxy_extras/migrations/20260823000000_add_spend_logs_api_key_starttime_index/migration.sql
@@ -1,2 +1,6 @@
--- CreateIndex
-CREATE INDEX IF NOT EXISTS "LiteLLM_SpendLogs_api_key_startTime_idx" ON "LiteLLM_SpendLogs"("api_key", "startTime");
+-- The (api_key, startTime) index on LiteLLM_SpendLogs is built after migrate deploy,
+-- through litellm_proxy_extras/request_log_indexes.py: concurrently on a plain table and
+-- per partition on a partitioned one. The migration job builds it; a serving proxy that
+-- ran the migrations itself builds it in the background once it serves. A migration
+-- cannot do either without blocking spend-log writes or failing on a partitioned table.
+SELECT 1;
diff --git a/litellm-proxy-extras/litellm_proxy_extras/migrations/20260831120001_spend_logs_litellm_call_id_index/migration.sql b/litellm-proxy-extras/litellm_proxy_extras/migrations/20260831120001_spend_logs_litellm_call_id_index/migration.sql
index 62ad5c42ba7..7eba7fc9b97 100644
--- a/litellm-proxy-extras/litellm_proxy_extras/migrations/20260831120001_spend_logs_litellm_call_id_index/migration.sql
+++ b/litellm-proxy-extras/litellm_proxy_extras/migrations/20260831120001_spend_logs_litellm_call_id_index/migration.sql
@@ -1,12 +1,6 @@
--- CreateIndex (CONCURRENTLY)
---
--- Disclaimer:
--- - CREATE INDEX CONCURRENTLY cannot run inside a transaction. This migration must stay a
--- single statement so Prisma Migrate on PostgreSQL can apply it outside a transaction.
--- - Builds are slower and use more I/O than a blocking CREATE INDEX; if the build is
--- interrupted, Postgres may leave an INVALID index that must be dropped and recreated.
--- - Do not edit this file after it has been applied to any database: Prisma checksums
--- migrations; add a new migration instead.
--- - Requires PostgreSQL that supports CONCURRENTLY with IF NOT EXISTS (use a new migration
--- without IF NOT EXISTS if you must support older versions).
-CREATE INDEX CONCURRENTLY IF NOT EXISTS "LiteLLM_SpendLogs_litellm_call_id_idx" ON "LiteLLM_SpendLogs"("litellm_call_id");
+-- The litellm_call_id index on LiteLLM_SpendLogs is built after migrate deploy, through
+-- litellm_proxy_extras/request_log_indexes.py: concurrently on a plain table and per
+-- partition on a partitioned one. The migration job builds it; a serving proxy that ran
+-- the migrations itself builds it in the background once it serves. Postgres refuses
+-- CREATE INDEX CONCURRENTLY on a partitioned parent, so this migration no longer runs it.
+SELECT 1;
diff --git a/litellm-proxy-extras/litellm_proxy_extras/request_log_indexes.py b/litellm-proxy-extras/litellm_proxy_extras/request_log_indexes.py
new file mode 100644
index 00000000000..6c31e8364a9
--- /dev/null
+++ b/litellm-proxy-extras/litellm_proxy_extras/request_log_indexes.py
@@ -0,0 +1,463 @@
+"""The request-log indexes built after `prisma migrate deploy` instead of by a migration:
+by the migration job, or by a serving proxy that ran the migrations itself (in the
+background, once it serves).
+
+A migration cannot build them: a plain `CREATE INDEX` blocks spend-log inserts for the
+whole build, and `CREATE INDEX CONCURRENTLY` is refused on a partitioned parent
+(db_scripts/partition_spend_logs.sql). `REQUEST_LOG_INDEXES` is the one list to extend;
+names match what Prisma derives from the `@@index` declarations in schema.prisma, so an
+index a database already has is recognized and never rebuilt.
+"""
+
+import hashlib
+import random
+import re
+import time
+from collections.abc import Callable
+from dataclasses import dataclass
+from typing import TYPE_CHECKING, Final
+
+from litellm_proxy_extras._logging import logger
+from litellm_proxy_extras.migration_lock import held_migration_lock
+
+if TYPE_CHECKING:
+ import psycopg
+ from psycopg import sql
+
+
+@dataclass(frozen=True, slots=True)
+class RequestLogIndex:
+ """One index the migration job owns: the table, the exact Prisma index name and the
+ column list as it would be written after `ON
`."""
+
+ table: str
+ name: str
+ definition: str
+
+ @property
+ def columns(self) -> tuple[str, ...]:
+ return tuple(re.findall(r'"([^"]+)"', self.definition))
+
+ def partition_index_name(self, partition: str) -> str:
+ """The child index name for one partition, built the way Postgres names the
+ children of a partitioned index, and kept within the 63 byte identifier limit."""
+ name: Final = f"{partition}_{self.name.removeprefix(f'{self.table}_')}"
+ if len(name.encode()) <= _IDENTIFIER_MAX_BYTES:
+ return name
+ digest: Final = hashlib.sha256(name.encode()).hexdigest()[:_DIGEST_LENGTH]
+ budget: Final = _IDENTIFIER_MAX_BYTES - _DIGEST_LENGTH - 1
+ kept: Final = next(name[:length] for length in range(len(name), 0, -1) if len(name[:length].encode()) <= budget)
+ return f"{kept}_{digest}"
+
+
+REQUEST_LOG_INDEXES: Final = (
+ RequestLogIndex("LiteLLM_SpendLogs", "LiteLLM_SpendLogs_api_key_startTime_idx", '("api_key", "startTime")'),
+ RequestLogIndex("LiteLLM_SpendLogs", "LiteLLM_SpendLogs_litellm_call_id_idx", '("litellm_call_id")'),
+)
+
+_IDENTIFIER_MAX_BYTES: Final = 63
+_DDL_LOCK_TIMEOUT: Final = "200ms"
+_DDL_LOCK_ATTEMPTS: Final = 10
+_DDL_RETRY_BASE_SECONDS: Final = 0.25
+_DDL_RETRY_MAX_SECONDS: Final = 8.0
+_LOCK_HANDOVER_SECONDS: Final = 2.0
+_DIGEST_LENGTH: Final = 8
+_CREATE_INDEX_STATEMENT: Final = re.compile(
+ r'^\s*CREATE\s+(?:UNIQUE\s+)?INDEX\s+(?:CONCURRENTLY\s+)?(?:IF\s+NOT\s+EXISTS\s+)?"(?P[^"]+)"\s+ON\b',
+ re.IGNORECASE,
+)
+_TABLE_KIND_SQL: Final = "SELECT c.relkind = 'p' AS partitioned FROM pg_class c WHERE c.oid = to_regclass(%s)"
+_CHILDREN_WITHOUT_THE_INDEX_SQL: Final = (
+ "SELECT child.relname AS name, n.nspname AS schema, child.relkind = 'p' AS partitioned "
+ "FROM pg_inherits i JOIN pg_class child ON child.oid = i.inhrelid "
+ "JOIN pg_namespace n ON n.oid = child.relnamespace "
+ "WHERE i.inhparent = to_regclass(%s) AND NOT EXISTS ("
+ "SELECT 1 FROM pg_inherits attached JOIN pg_index x ON x.indexrelid = attached.inhrelid "
+ "WHERE attached.inhparent = to_regclass(%s) AND x.indrelid = child.oid) "
+ "ORDER BY child.relname"
+)
+_EQUIVALENT_INDEXES_SQL: Final = (
+ "SELECT i.relname AS name, x.indisvalid AS valid "
+ "FROM pg_index x JOIN pg_class i ON i.oid = x.indexrelid JOIN pg_am am ON am.oid = i.relam "
+ "WHERE x.indrelid = to_regclass(%s) AND i.relname <> %s AND am.amname = 'btree' AND NOT x.indisunique "
+ "AND x.indexprs IS NULL AND x.indpred IS NULL AND x.indnkeyatts = x.indnatts "
+ "AND NOT EXISTS (SELECT 1 FROM unnest(x.indoption::int2[]) o WHERE o <> 0) "
+ "AND NOT EXISTS (SELECT 1 FROM unnest(x.indclass::oid[]) c JOIN pg_opclass oc ON oc.oid = c WHERE NOT oc.opcdefault) "
+ "AND NOT EXISTS (SELECT 1 FROM unnest(x.indcollation::oid[]) WITH ORDINALITY c(coll, ord) "
+ "JOIN unnest(x.indkey::int2[]) WITH ORDINALITY k(attnum, ord) ON k.ord = c.ord "
+ "JOIN pg_attribute a ON a.attrelid = x.indrelid AND a.attnum = k.attnum "
+ "WHERE c.coll <> 0 AND c.coll <> a.attcollation) "
+ "AND (SELECT array_agg(a.attname::text ORDER BY k.ord) FROM unnest(x.indkey::int2[]) WITH ORDINALITY k(attnum, ord) "
+ "JOIN pg_attribute a ON a.attrelid = x.indrelid AND a.attnum = k.attnum) = %s::text[] "
+ "AND NOT EXISTS (SELECT 1 FROM pg_inherits WHERE inhrelid = x.indexrelid) "
+ "ORDER BY x.indisvalid DESC, i.relname"
+)
+_INDEX_STATE_SQL: Final = (
+ 'SELECT x.indisvalid AS valid, t.relname AS "table" '
+ "FROM pg_index x JOIN pg_class t ON t.oid = x.indrelid WHERE x.indexrelid = to_regclass(%s)"
+)
+
+
+@dataclass(frozen=True, slots=True)
+class _Relation:
+ name: str
+ schema: str
+ partitioned: bool
+
+
+@dataclass(frozen=True, slots=True)
+class _IndexState:
+ valid: bool
+ table: str
+
+
+@dataclass(frozen=True, slots=True)
+class _EquivalentIndex:
+ name: str
+ valid: bool
+
+
+@dataclass(frozen=True, slots=True)
+class _TableKind:
+ partitioned: bool
+
+
+def filter_request_log_index_diff(diff_sql: str, indexes: tuple[RequestLogIndex, ...] = REQUEST_LOG_INDEXES) -> str:
+ """The `prisma migrate diff` script without the statements that create a migration-job-owned
+ index, which the schema declares and the migrations deliberately do not build."""
+ names: Final = frozenset(index.name for index in indexes)
+ statements: Final = diff_sql.split(";")
+ kept: Final = tuple(statement for statement in statements if not _creates_one_of(statement, names))
+ return ";".join(kept) if any(part.strip() for part in kept) else ""
+
+
+def _creates_one_of(statement: str, names: frozenset[str]) -> bool:
+ match: Final = _CREATE_INDEX_STATEMENT.match(_without_comments(statement))
+ return match is not None and match["index"] in names
+
+
+def _without_comments(statement: str) -> str:
+ return "\n".join(line for line in statement.splitlines() if not line.lstrip().startswith("--"))
+
+
+def _connect(database_url: str) -> "psycopg.Connection[tuple[object, ...]]":
+ import psycopg
+
+ return psycopg.connect(database_url, connect_timeout=10, autocommit=True)
+
+
+def ensure_request_log_indexes(
+ database_url: str,
+ schema: str,
+ indexes: tuple[RequestLogIndex, ...] = REQUEST_LOG_INDEXES,
+ connect: "Callable[[str], psycopg.Connection[tuple[object, ...]]]" = _connect,
+) -> bool:
+ """Build every listed index that is missing or invalid. Each build step runs under
+ the migration coordinator lock, held per statement so a resolver booting on another
+ replica gets in between partitions rather than waiting for the whole table. Any
+ failure is logged and left for the next index build; the result says whether
+ every index ended up valid. Never raises."""
+ import psycopg
+
+ try:
+ with connect(database_url) as connection:
+ connection.execute("SET statement_timeout = 0")
+ results: Final = tuple(_ensure_index(connection, schema, index) for index in indexes)
+ except psycopg.Error as exc:
+ logger.warning("Could not build the request-log indexes, leaving them for the next index build: %s", exc)
+ return False
+ if not all(results):
+ logger.warning("Some request-log indexes are not in place yet, leaving them for the next index build")
+ return False
+ logger.info("Request-log indexes are all in place")
+ return True
+
+
+def _under_migration_lock(connection: "psycopg.Connection[tuple[object, ...]]", step: Callable[[], bool]) -> bool:
+ with held_migration_lock(connection) as held:
+ if not held:
+ logger.info(
+ "Another process holds the migration lock, leaving the request-log indexes to the next index build"
+ )
+ return False
+ return step()
+
+
+def _with_bounded_lock(
+ connection: "psycopg.Connection[tuple[object, ...]]", step: Callable[[], bool], what: str
+) -> bool:
+ """Run `step` under the migration lock with a short lock_timeout, so a DDL statement that has to wait for open
+ transactions holds new writes back for at most that long; retry with capped exponential backoff, holding the
+ migration lock per attempt only and releasing it while sleeping. False when another process holds the migration
+ lock or every attempt timed out."""
+ import psycopg
+ from psycopg import sql
+
+ for attempt in range(_DDL_LOCK_ATTEMPTS):
+ if attempt:
+ time.sleep(min(_DDL_RETRY_MAX_SECONDS, _DDL_RETRY_BASE_SECONDS * 2.0**attempt) * random.uniform(0.5, 1.0))
+ connection.execute(sql.SQL("SET lock_timeout = {}").format(sql.Literal(_DDL_LOCK_TIMEOUT)))
+ try:
+ return _under_migration_lock(connection, step)
+ except psycopg.errors.LockNotAvailable:
+ logger.info("Waiting for open transactions before %s", what)
+ finally:
+ connection.execute("SET lock_timeout = 0")
+ logger.warning(
+ "Could not get the lock for %s without holding writes back, leaving it for the next index build", what
+ )
+ return False
+
+
+def _ensure_index(connection: "psycopg.Connection[tuple[object, ...]]", schema: str, index: RequestLogIndex) -> bool:
+ from psycopg.rows import class_row
+
+ with connection.cursor(row_factory=class_row(_TableKind)) as cursor:
+ table: Final = cursor.execute(_TABLE_KIND_SQL, (_regclass_name(connection, schema, index.table),)).fetchone()
+ if table is None:
+ logger.info("Table %s does not exist yet, skipping index %s", index.table, index.name)
+ return True
+ if table.partitioned:
+ return build_index_on_partitioned_table(connection, schema, index)
+ return _build_leaf_index(connection, schema, index.table, index.name, index)
+
+
+def _regclass_name(connection: "psycopg.Connection[tuple[object, ...]]", schema: str, name: str) -> str:
+ from psycopg import sql
+
+ return sql.Identifier(schema, name).as_string(connection)
+
+
+def _create_index_statement(
+ connection: "psycopg.Connection[tuple[object, ...]]", prefix: "sql.Composed", definition: str
+) -> bytes:
+ return (prefix.as_string(connection) + definition).encode()
+
+
+def _index_state(connection: "psycopg.Connection[tuple[object, ...]]", schema: str, index: str) -> "_IndexState | None":
+ from psycopg.rows import class_row
+
+ with connection.cursor(row_factory=class_row(_IndexState)) as cursor:
+ return cursor.execute(_INDEX_STATE_SQL, (_regclass_name(connection, schema, index),)).fetchone()
+
+
+def _equivalent_indexes(
+ connection: "psycopg.Connection[tuple[object, ...]]",
+ schema: str,
+ table: str,
+ name: str,
+ index: RequestLogIndex,
+) -> tuple[_EquivalentIndex, ...]:
+ """The indexes on `table` other than `name` with the same definition: default btree
+ over the same columns in the same order, no expression, predicate, DESC or custom
+ opclass or collation, and not attached under a partitioned index. Valid ones first."""
+ from psycopg.rows import class_row
+
+ with connection.cursor(row_factory=class_row(_EquivalentIndex)) as cursor:
+ return tuple(
+ cursor.execute(
+ _EQUIVALENT_INDEXES_SQL, (_regclass_name(connection, schema, table), name, list(index.columns))
+ ).fetchall()
+ )
+
+
+def _adopt_equivalent_index(
+ connection: "psycopg.Connection[tuple[object, ...]]",
+ schema: str,
+ table: str,
+ name: str,
+ index: RequestLogIndex,
+) -> bool:
+ """Rename a valid index of the same definition under another name (an operator's
+ hand-built copy, say) to the name this code expects, instead of building a second
+ one. RENAME on an index is a catalog change that lets writes through."""
+ from psycopg import sql
+
+ equivalent: Final = next(
+ (found for found in _equivalent_indexes(connection, schema, table, name, index) if found.valid), None
+ )
+ if equivalent is None:
+ return False
+ logger.info(
+ "Renaming the equivalent index %s on %s to %s instead of building a second one", equivalent.name, table, name
+ )
+ connection.execute(
+ sql.SQL("ALTER INDEX {} RENAME TO {}").format(sql.Identifier(schema, equivalent.name), sql.Identifier(name))
+ )
+ return True
+
+
+def _report_second_copies(
+ connection: "psycopg.Connection[tuple[object, ...]]",
+ schema: str,
+ table: str,
+ name: str,
+ index: RequestLogIndex,
+ concurrently: bool,
+) -> None:
+ """Log every other index of the same definition with the statement that removes it.
+ Dropping is the operator's call: a second copy costs writes and disk, never results."""
+ from psycopg import sql
+
+ drop: Final = "DROP INDEX CONCURRENTLY" if concurrently else "DROP INDEX"
+ for copy in _equivalent_indexes(connection, schema, table, name, index):
+ logger.warning(
+ "Index %s on %s is a second copy of %s and only costs writes and disk; remove it with: %s %s",
+ copy.name,
+ table,
+ name,
+ drop,
+ sql.Identifier(schema, copy.name).as_string(connection),
+ )
+
+
+def _children_without_the_index(
+ connection: "psycopg.Connection[tuple[object, ...]]", schema: str, table: str, index: str
+) -> tuple[_Relation, ...]:
+ from psycopg.rows import class_row
+
+ with connection.cursor(row_factory=class_row(_Relation)) as cursor:
+ return tuple(
+ cursor.execute(
+ _CHILDREN_WITHOUT_THE_INDEX_SQL,
+ (_regclass_name(connection, schema, table), _regclass_name(connection, schema, index)),
+ ).fetchall()
+ )
+
+
+def _build_leaf_index(
+ connection: "psycopg.Connection[tuple[object, ...]]",
+ schema: str,
+ table: str,
+ name: str,
+ index: RequestLogIndex,
+) -> bool:
+ """Build one plain table's or partition's index with CONCURRENTLY so writes keep
+ flowing. The catalog is read under the migration lock, so a replica that saw an
+ invalid index before the lock finds the valid one another replica just built and
+ leaves it. An invalid index left by an interrupted build is dropped and rebuilt; a
+ valid index of the same definition under another name is renamed rather than
+ duplicated; an index of that name on another table is a collision this code will
+ not touch."""
+ from psycopg import sql
+
+ def build() -> bool:
+ existing: Final = _index_state(connection, schema, name)
+ if existing is not None and existing.table != table:
+ logger.warning(
+ "Index %s already exists on %s rather than %s, leaving it alone", name, existing.table, table
+ )
+ return False
+ if existing is not None and existing.valid:
+ return True
+ if existing is not None:
+ logger.info("Dropping the invalid index %s left by an interrupted build on %s", name, table)
+ connection.execute(sql.SQL("DROP INDEX CONCURRENTLY {}").format(sql.Identifier(schema, name)))
+ elif _adopt_equivalent_index(connection, schema, table, name, index):
+ return True
+ logger.info("Building index %s on %s concurrently", name, table)
+ prefix: Final = sql.SQL("CREATE INDEX CONCURRENTLY IF NOT EXISTS {} ON {} ").format(
+ sql.Identifier(name), sql.Identifier(schema, table)
+ )
+ connection.execute(_create_index_statement(connection, prefix, index.definition))
+ built: Final = _index_state(connection, schema, name)
+ return built is not None and built.valid
+
+ current: Final = _index_state(connection, schema, name)
+ if current is None or not current.valid or current.table != table:
+ if not _under_migration_lock(connection, build):
+ return False
+ time.sleep(_LOCK_HANDOVER_SECONDS)
+ _report_second_copies(connection, schema, table, name, index, concurrently=True)
+ return True
+
+
+def build_index_on_partitioned_table(
+ connection: "psycopg.Connection[tuple[object, ...]]",
+ schema: str,
+ index: RequestLogIndex,
+ table: "str | None" = None,
+ name: "str | None" = None,
+) -> bool:
+ """Build the index the way Postgres allows on a partitioned parent: a metadata-only
+ parent index ON ONLY the parent, one CONCURRENTLY build per partition, and ATTACH
+ PARTITION for each child. Partitions that are themselves partitioned get the same
+ treatment one level down. Every step checks the catalog before acting, so an
+ interrupted run resumes where it stopped and a second run finds nothing to do; a
+ parent or child index of the same definition under another name is renamed and
+ used rather than duplicated. The connection must be in autocommit mode. True when
+ the parent index ends up valid."""
+
+ parent_table: Final = index.table if table is None else table
+ parent_index: Final = index.name if name is None else name
+ existing: Final = _index_state(connection, schema, parent_index)
+ if existing is not None and existing.table != parent_table:
+ logger.warning(
+ "Index %s already exists on %s rather than %s, leaving it alone", parent_index, existing.table, parent_table
+ )
+ return False
+ if existing is None and not _with_bounded_lock(
+ connection,
+ lambda: (
+ _adopt_equivalent_index(connection, schema, parent_table, parent_index, index)
+ or _create_parent_index(connection, schema, parent_index, parent_table, index)
+ ),
+ f"creating the parent index {parent_index}",
+ ):
+ return False
+ children: Final = _children_without_the_index(connection, schema, parent_table, parent_index)
+ if not all(_attach_child_index(connection, schema, parent_index, child, index) for child in children):
+ return False
+ final: Final = _index_state(connection, schema, parent_index)
+ if final is None or not final.valid:
+ return False
+ _report_second_copies(connection, schema, parent_table, parent_index, index, concurrently=False)
+ return True
+
+
+def _create_parent_index(
+ connection: "psycopg.Connection[tuple[object, ...]]",
+ schema: str,
+ name: str,
+ table: str,
+ index: RequestLogIndex,
+) -> bool:
+ """Create the metadata-only parent index. The caller bounds Postgres's SHARE lock wait on the parent."""
+ from psycopg import sql
+
+ prefix: Final = sql.SQL("CREATE INDEX IF NOT EXISTS {} ON ONLY {} ").format(
+ sql.Identifier(name), sql.Identifier(schema, table)
+ )
+ statement: Final = _create_index_statement(connection, prefix, index.definition)
+ connection.execute(statement)
+ return True
+
+
+def _attach_child_index(
+ connection: "psycopg.Connection[tuple[object, ...]]",
+ schema: str,
+ parent_index: str,
+ child: _Relation,
+ index: RequestLogIndex,
+) -> bool:
+ from psycopg import sql
+
+ child_index: Final = index.partition_index_name(child.name)
+ built: Final = (
+ build_index_on_partitioned_table(connection, child.schema, index, child.name, child_index)
+ if child.partitioned
+ else _build_leaf_index(connection, child.schema, child.name, child_index, index)
+ )
+ if not built:
+ return False
+
+ def attach() -> bool:
+ connection.execute(
+ sql.SQL("ALTER INDEX {} ATTACH PARTITION {}").format(
+ sql.Identifier(schema, parent_index), sql.Identifier(child.schema, child_index)
+ )
+ )
+ logger.info("Attached index %s on partition %s to %s", child_index, child.name, parent_index)
+ return True
+
+ return _with_bounded_lock(connection, attach, f"attaching {child_index}")
diff --git a/litellm-proxy-extras/litellm_proxy_extras/utils.py b/litellm-proxy-extras/litellm_proxy_extras/utils.py
index 2f74df63367..3acc19d397d 100644
--- a/litellm-proxy-extras/litellm_proxy_extras/utils.py
+++ b/litellm-proxy-extras/litellm_proxy_extras/utils.py
@@ -5,6 +5,7 @@ import re
import shutil
import subprocess
import tempfile
+import threading
import time
from collections.abc import Callable
from dataclasses import dataclass, replace
@@ -13,6 +14,7 @@ from typing import TYPE_CHECKING, Final, Optional
from litellm_proxy_extras import prisma_toolchain
from litellm_proxy_extras._logging import logger
+from litellm_proxy_extras.migration_lock import held_migration_lock
from litellm_proxy_extras.prisma_toolchain import (
PRISMA_COMMAND_TIMEOUT_ENV_VAR,
PRISMA_MIGRATE_DEPLOY_TIMEOUT_ENV_VAR,
@@ -24,6 +26,7 @@ from litellm_proxy_extras.replica_identity import (
REPLICA_IDENTITY_FULL_ENV_VAR,
apply_replica_identity_full,
)
+from litellm_proxy_extras.request_log_indexes import ensure_request_log_indexes, filter_request_log_index_diff
if TYPE_CHECKING:
import psycopg
@@ -433,6 +436,21 @@ class ProxyExtrasDBManager:
return True
return False
+ @staticmethod
+ def _filter_migration_job_owned_drift(diff_sql: str, partitioned: bool | None = None) -> str:
+ """The drift script without the indexes the migration job builds (the schema
+ declares them, the migrations deliberately do not) and, when LiteLLM_SpendLogs
+ is partitioned, without its primary-key rewrite and partitioning artifacts."""
+ without_indexes: Final = filter_request_log_index_diff(diff_sql)
+ is_partitioned: Final = ProxyExtrasDBManager.spend_logs_is_partitioned() if partitioned is None else partitioned
+ if not is_partitioned:
+ return without_indexes
+ logger.info(
+ "LiteLLM_SpendLogs is partitioned; removed its primary-key "
+ "rewrite and partitioning artifacts from the drift script"
+ )
+ return filter_partitioned_spend_logs_diff(without_indexes)
+
@staticmethod
def _resolve_all_migrations(
migrations_dir: str, schema_path: str, mark_all_applied: bool = True
@@ -513,21 +531,14 @@ class ProxyExtrasDBManager:
return
logger.info(f"Migration diff created at {diff_sql_path}")
- if ProxyExtrasDBManager.spend_logs_is_partitioned():
- filtered_sql = filter_partitioned_spend_logs_diff(
- diff_sql_path.read_text()
- )
- diff_sql_path.write_text(filtered_sql)
- logger.info(
- "LiteLLM_SpendLogs is partitioned; removed its primary-key "
- "rewrite and partitioning artifacts from the drift script"
- )
- if not filtered_sql.strip():
- logger.info("Drift script is empty after filtering; nothing to apply")
- if not mark_all_applied:
- return
- ProxyExtrasDBManager._mark_migrations_applied(migrations_dir)
+ filtered_sql: Final = ProxyExtrasDBManager._filter_migration_job_owned_drift(diff_sql_path.read_text())
+ diff_sql_path.write_text(filtered_sql)
+ if not filtered_sql.strip():
+ logger.info("Drift script is empty after filtering; nothing to apply")
+ if not mark_all_applied:
return
+ ProxyExtrasDBManager._mark_migrations_applied(migrations_dir)
+ return
# 2. Run prisma db execute to apply the migration
applied_ok = False
@@ -800,7 +811,7 @@ class ProxyExtrasDBManager:
conn.execute(statement)
except psycopg.Error as e:
logger.warning(
- "Could not repair invalid index %s.%s, will retry on the next startup. "
+ "Could not repair invalid index %s.%s, will retry on the next database setup run. "
"If this keeps happening, run `%s` by hand as the index owner. Error: %s",
index.schema,
index.name,
@@ -811,16 +822,21 @@ class ProxyExtrasDBManager:
logger.info("%s invalid index %s.%s", action, index.schema, index.name)
@staticmethod
- def repair_invalid_indexes(lock_timeout: str = "30s") -> bool:
+ def repair_invalid_indexes(
+ lock_timeout: str = "30s",
+ repair: "Callable[[psycopg.Connection[tuple[str, str, str]], _InvalidIndex], None] | None" = None,
+ ) -> bool:
"""Rebuild LiteLLM indexes an interrupted CREATE INDEX CONCURRENTLY left
INVALID (a migration deadlock between replicas is the usual cause; the
retried migration skips them because of IF NOT EXISTS). Never raises:
returns True when no invalid index remains, False when the repair was
- skipped or failed and will be retried on the next startup. Looks in the
+ skipped or failed and will be retried on the next database setup run. Looks in the
schema DATABASE_URL names, the only URL Prisma migrates through, but
connects over DIRECT_URL when set: the session settings, the advisory
lock and REINDEX CONCURRENTLY all need one server session, which a
- transaction pooler does not give."""
+ transaction pooler does not give. Each rebuild holds the migration
+ coordinator lock on its own, like the migration job's index build, so a resolver
+ booting on another replica waits for one index at most."""
prisma_url: Final = os.getenv("DATABASE_URL")
if not prisma_url:
return False
@@ -856,20 +872,53 @@ class ProxyExtrasDBManager:
if lock_row is None or not lock_row[0]:
logger.info("Another replica is already rebuilding the invalid indexes, skipping")
return False
- for index in ProxyExtrasDBManager._invalid_litellm_indexes(conn, schema):
- ProxyExtrasDBManager._repair_index(conn, index)
+ repair_one: Final = repair or ProxyExtrasDBManager._repair_index
+ repaired: Final = all(
+ ProxyExtrasDBManager._repair_under_migration_lock(conn, schema, index, repair_one)
+ for index in found
+ )
+ if not repaired:
+ return False
remaining: Final = ProxyExtrasDBManager._invalid_litellm_indexes(conn, schema)
except psycopg.Error as e:
- logger.warning("Could not check for invalid indexes, will retry on the next startup. Error: %s", e)
+ logger.warning(
+ "Could not check for invalid indexes, will retry on the next database setup run. Error: %s", e
+ )
return False
return not remaining
+ @staticmethod
+ def _repair_under_migration_lock(
+ conn: "psycopg.Connection[tuple[str, str, str]]",
+ schema: str,
+ index: _InvalidIndex,
+ repair: "Callable[[psycopg.Connection[tuple[str, str, str]], _InvalidIndex], None]",
+ ) -> bool:
+ """Rebuild one index under the migration coordinator lock, skipping it when a
+ migration job finished or dropped it in the meantime. False when another process
+ holds the lock, so the check waits for the next database setup run."""
+ with held_migration_lock(conn) as held:
+ if not held:
+ logger.info(
+ "Another process is building indexes under the migration lock, leaving the "
+ "invalid index check to the next database setup run"
+ )
+ return False
+ still_invalid: Final = ProxyExtrasDBManager._invalid_litellm_indexes(conn, schema)
+ if any(found.schema == index.schema and found.name == index.name for found in still_invalid):
+ repair(conn, index)
+ return True
+
@staticmethod
def _setup_database_v2(use_migrate: bool) -> bool:
if not use_migrate:
return ProxyExtrasDBManager._run_database_v2(False)
from litellm_proxy_extras.migration_lock import migration_environment, migration_lock
- from litellm_proxy_extras.migration_recovery import baseline_current_schema, recover_completed_migration
+ from litellm_proxy_extras.migration_recovery import (
+ baseline_current_schema,
+ recover_completed_migration,
+ roll_back_failed_inert_migration,
+ )
database_url: Final = os.environ.get("DATABASE_URL")
if not database_url:
@@ -884,7 +933,9 @@ class ProxyExtrasDBManager:
if not migration.is_file():
return False
with migration_lock(lock_url) as coordinator:
- return recover_completed_migration(coordinator, schema, migration)
+ return recover_completed_migration(coordinator, schema, migration) or roll_back_failed_inert_migration(
+ coordinator, schema, migration
+ )
def baseline_existing(migrations_dir: str) -> None:
with migration_lock(lock_url) as coordinator:
@@ -1177,13 +1228,16 @@ class ProxyExtrasDBManager:
)
@staticmethod
- def setup_database(
- use_migrate: bool = False, use_v2_resolver: bool = False
- ) -> bool:
+ def setup_database(use_migrate: bool = False, use_v2_resolver: bool = False) -> bool:
"""
Set up the database using either prisma migrate or prisma db push
Uses migrations from litellm-proxy-extras package
+ The request-log indexes in `REQUEST_LOG_INDEXES` are not built here: the
+ migration job builds them through `run_migration_job`, and a serving proxy that
+ ran the migrations itself starts them through `start_request_log_index_build`
+ once it is ready to serve.
+
Args:
use_migrate: Whether to use prisma migrate instead of db push
use_v2_resolver: Opt into the v2 migration resolver (safer during
@@ -1200,10 +1254,48 @@ class ProxyExtrasDBManager:
migrated = ProxyExtrasDBManager._run_migrations(
use_migrate=use_migrate, use_v2_resolver=use_v2_resolver
)
- if migrated:
- ProxyExtrasDBManager.repair_invalid_indexes()
- ProxyExtrasDBManager.apply_replica_identity_full_if_requested()
- return migrated
+ if not migrated:
+ return False
+ ProxyExtrasDBManager.repair_invalid_indexes()
+ ProxyExtrasDBManager.apply_replica_identity_full_if_requested()
+ return True
+
+ @staticmethod
+ def build_request_log_indexes(build: Callable[[str, str], bool] = ensure_request_log_indexes) -> bool:
+ """Build the indexes in `REQUEST_LOG_INDEXES` on the writer, in the schema the
+ migrations target. Idempotent and never raises; False when an index is still
+ missing or invalid, so the migration job reports it and gets rerun instead of
+ leaving the table unindexed until the next deploy."""
+ database_url: Final = os.environ.get("DATABASE_URL")
+ if not database_url:
+ return True
+ direct_url: Final = ProxyExtrasDBManager._strip_prisma_query_params(
+ os.environ.get("DIRECT_URL") or database_url
+ )
+ schema: Final = ProxyExtrasDBManager._prisma_schema_param(database_url) or "public"
+ return build(direct_url, schema)
+
+ @staticmethod
+ def run_migration_job(
+ use_migrate: bool = False,
+ use_v2_resolver: bool = False,
+ setup: Callable[[bool, bool], bool] = setup_database,
+ build: Callable[[], bool] = build_request_log_indexes,
+ ) -> bool:
+ """The migration job's whole run: `setup_database`, then the request-log indexes,
+ built synchronously so the job exits only once they are in place. False when the
+ migrations failed or an index could not be built, so the Job is rerun."""
+ return setup(use_migrate, use_v2_resolver) and build()
+
+ @staticmethod
+ def start_request_log_index_build(build: Callable[[], bool] = build_request_log_indexes) -> threading.Thread:
+ """A serving proxy that ran the migrations itself (schema updates not disabled)
+ builds the request-log indexes on a daemon thread, so a long build never delays
+ readiness. A build that could not finish is logged and picked up by the next boot
+ or the migration job."""
+ thread: Final = threading.Thread(target=build, name="litellm-request-log-indexes", daemon=True)
+ thread.start()
+ return thread
@staticmethod
def _run_migrations(use_migrate: bool, use_v2_resolver: bool) -> bool:
@@ -1247,15 +1339,16 @@ class ProxyExtrasDBManager:
logger.info("✅ Post-migration sanity check completed")
return True
except subprocess.CalledProcessError as e:
- logger.info(f"prisma db error: {e.stderr}, e: {e.stdout}")
- if "P3009" in e.stderr:
+ stderr: Final = str(e.stderr or "")
+ logger.info(f"prisma db error: {stderr}, e: {e.stdout}")
+ if "P3009" in stderr:
# Extract the failed migration name from the error message
migration_match = re.search(
- r"`(\d+_.*)` migration", e.stderr
+ r"`(\d+_.*)` migration", stderr
)
if migration_match:
failed_migration = migration_match.group(1)
- if ProxyExtrasDBManager._is_idempotent_error(e.stderr):
+ if ProxyExtrasDBManager._is_idempotent_error(stderr):
logger.info(
f"Migration {failed_migration} failed due to idempotent error (e.g., column already exists), resolving as applied"
)
@@ -1311,8 +1404,8 @@ class ProxyExtrasDBManager:
f"✅ Migration {failed_migration} marked as rolled back... retrying"
)
elif (
- "P3005" in e.stderr
- and "database schema is not empty" in e.stderr
+ "P3005" in stderr
+ and "database schema is not empty" in stderr
):
logger.info(
"Database schema is not empty, creating baseline migration. In read-only file system, please set an environment variable `LITELLM_MIGRATION_DIR` to a writable directory to enable migrations. Learn more - https://docs.litellm.ai/docs/proxy/prod#read-only-file-system"
@@ -1326,13 +1419,13 @@ class ProxyExtrasDBManager:
)
logger.info("✅ All migrations resolved.")
return True
- elif "P3018" in e.stderr:
+ elif "P3018" in stderr:
# Check if this is a permission error or idempotent error
- if ProxyExtrasDBManager._is_permission_error(e.stderr):
+ if ProxyExtrasDBManager._is_permission_error(stderr):
# Permission errors should NOT be marked as applied
# Extract migration name for logging
migration_match = re.search(
- r"Migration name: (\d+_.*)", e.stderr
+ r"Migration name: (\d+_.*)", stderr
)
migration_name = (
migration_match.group(1)
@@ -1342,7 +1435,7 @@ class ProxyExtrasDBManager:
logger.error(
f"❌ Migration {migration_name} failed due to insufficient permissions. "
- f"Please check database user privileges. Error: {e.stderr}"
+ f"Please check database user privileges. Error: {stderr}"
)
# Mark as rolled back and exit with error
@@ -1365,7 +1458,7 @@ class ProxyExtrasDBManager:
f"was NOT applied. Please grant necessary database permissions and retry."
) from e
- elif ProxyExtrasDBManager._is_idempotent_error(e.stderr):
+ elif ProxyExtrasDBManager._is_idempotent_error(stderr):
# Idempotent errors mean the migration has effectively been applied
logger.info(
"Migration failed due to idempotent error (e.g., column already exists), "
@@ -1373,7 +1466,7 @@ class ProxyExtrasDBManager:
)
# Extract the migration name from the error message
migration_match = re.search(
- r"Migration name: (\d+_.*)", e.stderr
+ r"Migration name: (\d+_.*)", stderr
)
if migration_match:
migration_name = migration_match.group(1)
@@ -1422,7 +1515,7 @@ class ProxyExtrasDBManager:
logger.warning(
f"P3018 error encountered but could not classify "
f"as permission or idempotent error. "
- f"Error: {e.stderr}"
+ f"Error: {stderr}"
)
raise
else:
diff --git a/litellm-proxy-extras/pyproject.toml b/litellm-proxy-extras/pyproject.toml
index 2e2f3f2ce5a..c6d6060acfa 100644
--- a/litellm-proxy-extras/pyproject.toml
+++ b/litellm-proxy-extras/pyproject.toml
@@ -1,6 +1,6 @@
[project]
name = "litellm-proxy-extras"
-version = "0.4.103"
+version = "0.4.104"
description = "Additional files for the LiteLLM Proxy. Reduces the size of the main litellm package."
readme = "README.md"
requires-python = ">=3.9"
@@ -30,7 +30,7 @@ required-version = ">=0.10.9"
module-root = ""
[tool.commitizen]
-version = "0.4.103"
+version = "0.4.104"
version_files = [
"pyproject.toml:^version",
"../pyproject.toml:litellm-proxy-extras==",
diff --git a/litellm-rust/Cargo.lock b/litellm-rust/Cargo.lock
index ff0eafee47e..40552b19e43 100644
--- a/litellm-rust/Cargo.lock
+++ b/litellm-rust/Cargo.lock
@@ -97,6 +97,53 @@ version = "1.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "03918c3dbd7701a85c6b9887732e2921175f26c350b4563841d0958c21d57e6d"
+[[package]]
+name = "askama"
+version = "0.16.1"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "6024d73179f43f15ccd2b881bfea6fee7f3a46ec53f33b52210dea749ebebaa4"
+dependencies = [
+ "askama_macros",
+ "itoa",
+ "percent-encoding",
+ "serde",
+ "serde_json",
+]
+
+[[package]]
+name = "askama_derive"
+version = "0.16.1"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "071ee5ebf2138e3ad180e0aacf6940c2cab5e6d8333741d9925c7bee2b153f39"
+dependencies = [
+ "askama_parser",
+ "memchr",
+ "proc-macro2",
+ "quote",
+ "rustc-hash",
+ "syn 3.0.6",
+]
+
+[[package]]
+name = "askama_macros"
+version = "0.16.1"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "643e1c7cbb6aec1d920332fe51a7c0d8219e273dcb8602db03f5263e4d16487b"
+dependencies = [
+ "askama_derive",
+]
+
+[[package]]
+name = "askama_parser"
+version = "0.16.1"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "2c5ae75772275d268b03ab8bdccdd12117b6169ee23256942b34e46c9f476583"
+dependencies = [
+ "rustc-hash",
+ "unicode-ident",
+ "winnow 1.0.4",
+]
+
[[package]]
name = "asn1-rs"
version = "0.7.2"
@@ -4038,6 +4085,26 @@ dependencies = [
"strum",
]
+[[package]]
+name = "litellm-migrate"
+version = "0.1.0"
+dependencies = [
+ "litellm-migrate-macros",
+ "rstest",
+]
+
+[[package]]
+name = "litellm-migrate-macros"
+version = "0.1.0"
+dependencies = [
+ "proc-macro2",
+ "quote",
+ "rstest",
+ "syn 2.0.119",
+ "tempfile",
+ "thiserror 2.0.19",
+]
+
[[package]]
name = "litellm-model-catalog"
version = "0.1.0"
@@ -4384,11 +4451,17 @@ dependencies = [
name = "litellm-traces"
version = "0.1.0"
dependencies = [
+ "askama",
"base64 0.22.1",
"criterion",
"flate2",
+ "futures-util",
+ "hmac 0.12.1",
+ "indexmap 2.14.0",
"litellm-http",
+ "litellm-migrate",
"litellm-storage-clickhouse",
+ "moka",
"opentelemetry-proto",
"prost",
"rstest",
@@ -4400,6 +4473,7 @@ dependencies = [
"thiserror 2.0.19",
"time",
"tokio",
+ "url",
"wiremock",
]
diff --git a/litellm-rust/Cargo.toml b/litellm-rust/Cargo.toml
index 8d837c2d31b..f4cb2ecbe59 100644
--- a/litellm-rust/Cargo.toml
+++ b/litellm-rust/Cargo.toml
@@ -14,6 +14,8 @@ litellm-router = { path = "crates/router" }
litellm-tracing = { path = "crates/tracing" }
litellm-traces = { path = "crates/traces" }
litellm-storage-clickhouse = { path = "crates/storage-clickhouse" }
+litellm-migrate = { path = "crates/migrate" }
+litellm-migrate-macros = { path = "crates/migrate-macros" }
litellm-core = { path = "crates/core" }
litellm-gateway-mcp = { path = "crates/gateway-mcp" }
litellm-gateway = { path = "crates/gateway" }
@@ -63,6 +65,7 @@ litellm-token-counter-tiktoken = { path = "crates/token-counter-tiktoken" }
litellm-host-python = { path = "crates/host-python" }
litellm-python-compat = { path = "crates/python-compat" }
+askama = { version = "0.16.1", default-features = false, features = ["derive", "std"] }
tracing = "0.1"
axum = { version = "0.8.9", default-features = false, features = ["http1", "tokio", "multipart"] }
axum-login = "0.18.0"
@@ -93,7 +96,10 @@ 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"
+syn = { version = "2", default-features = false }
sqlx = { version = "0.9.0", default-features = false, features = ["json", "macros", "postgres", "runtime-tokio", "chrono", "tls-rustls-ring-native-roots"] }
+proc-macro2 = "1"
+quote = "1"
subtle = "2"
thiserror = "2.0"
tokenizers = { version = "0.23.1", default-features = false, features = ["onig"] }
diff --git a/litellm-rust/crates/migrate-macros/Cargo.toml b/litellm-rust/crates/migrate-macros/Cargo.toml
new file mode 100644
index 00000000000..5cd68415ca2
--- /dev/null
+++ b/litellm-rust/crates/migrate-macros/Cargo.toml
@@ -0,0 +1,19 @@
+[package]
+name = "litellm-migrate-macros"
+version = "0.1.0"
+edition.workspace = true
+license.workspace = true
+repository.workspace = true
+
+[lib]
+proc-macro = true
+
+[dependencies]
+proc-macro2.workspace = true
+quote.workspace = true
+syn = { workspace = true, features = ["parsing", "printing", "proc-macro"] }
+thiserror.workspace = true
+
+[dev-dependencies]
+rstest.workspace = true
+tempfile.workspace = true
diff --git a/litellm-rust/crates/migrate-macros/src/error.rs b/litellm-rust/crates/migrate-macros/src/error.rs
new file mode 100644
index 00000000000..9833009517b
--- /dev/null
+++ b/litellm-rust/crates/migrate-macros/src/error.rs
@@ -0,0 +1,21 @@
+use std::io;
+
+#[derive(Debug, thiserror::Error)]
+pub enum Error {
+ #[error("could not read migrations directory `{path}`")]
+ ReadDirectory {
+ path: String,
+ #[source]
+ source: io::Error,
+ },
+ #[error(
+ "migration name `{name}` must be `_.sql` with a `[a-z0-9_]` description"
+ )]
+ InvalidName { name: String },
+ #[error("migration version `{version}` is declared more than once")]
+ DuplicateVersion { version: u64 },
+ #[error("migrations directory `{path}` contains no migrations")]
+ Empty { path: String },
+ #[error("migration path `{path}` is not valid UTF-8")]
+ NonUtf8Path { path: String },
+}
diff --git a/litellm-rust/crates/migrate-macros/src/lib.rs b/litellm-rust/crates/migrate-macros/src/lib.rs
new file mode 100644
index 00000000000..501f59e6fc2
--- /dev/null
+++ b/litellm-rust/crates/migrate-macros/src/lib.rs
@@ -0,0 +1,199 @@
+mod error;
+
+use std::path::{Path, PathBuf};
+
+use error::Error;
+use proc_macro::TokenStream;
+use quote::quote;
+use syn::LitStr;
+
+struct Entry {
+ version: u64,
+ description: String,
+ path: PathBuf,
+}
+
+fn resolve(dir: &Path) -> Result, Error> {
+ let mut entries = Vec::new();
+ let files = std::fs::read_dir(dir).map_err(|source| Error::ReadDirectory {
+ path: dir.display().to_string(),
+ source,
+ })?;
+ for file in files {
+ let file = file.map_err(|source| Error::ReadDirectory {
+ path: dir.display().to_string(),
+ source,
+ })?;
+ let path = file.path();
+ let name = path
+ .file_name()
+ .and_then(|name| name.to_str())
+ .ok_or_else(|| Error::NonUtf8Path {
+ path: path.display().to_string(),
+ })?
+ .to_owned();
+ let invalid = || Error::InvalidName { name: name.clone() };
+ let stem = name
+ .strip_suffix(".sql")
+ .filter(|_| file.file_type().is_ok_and(|kind| kind.is_file()))
+ .and_then(|stem| stem.split_once('_'))
+ .filter(|(version, description)| {
+ !version.is_empty()
+ && version.bytes().all(|b| b.is_ascii_digit())
+ && !description.is_empty()
+ && description
+ .bytes()
+ .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'_')
+ })
+ .ok_or_else(invalid)?;
+ let version = stem.0.parse::().map_err(|_| invalid())?;
+ entries.push(Entry {
+ version,
+ description: stem.1.to_owned(),
+ path,
+ });
+ }
+ if entries.is_empty() {
+ return Err(Error::Empty {
+ path: dir.display().to_string(),
+ });
+ }
+ entries.sort_by_key(|entry| entry.version);
+ for pair in entries.windows(2) {
+ if pair[0].version == pair[1].version {
+ return Err(Error::DuplicateVersion {
+ version: pair[0].version,
+ });
+ }
+ }
+ Ok(entries)
+}
+
+fn resolve_input(lit: &LitStr) -> Result, Error> {
+ let root = std::env::var("CARGO_MANIFEST_DIR")
+ .map(PathBuf::from)
+ .unwrap_or_default();
+ let dir = root.join(lit.value());
+ let dir = dir.canonicalize().map_err(|source| Error::ReadDirectory {
+ path: dir.display().to_string(),
+ source,
+ })?;
+ if dir.to_str().is_none() {
+ return Err(Error::NonUtf8Path {
+ path: dir.display().to_string(),
+ });
+ }
+ resolve(&dir)
+}
+
+#[proc_macro]
+pub fn migrate(input: TokenStream) -> TokenStream {
+ let lit = syn::parse_macro_input!(input as LitStr);
+ match resolve_input(&lit) {
+ Ok(entries) => {
+ let migrations = entries.iter().map(|entry| {
+ let version = entry.version;
+ let description = &entry.description;
+ let path = entry
+ .path
+ .to_str()
+ .expect("canonical migration path is UTF-8");
+ quote! {
+ ::litellm_migrate::Migration {
+ version: #version,
+ description: #description,
+ sql: ::core::include_str!(#path),
+ }
+ }
+ });
+ quote! { &[#(#migrations),*] }.into()
+ }
+ Err(err) => syn::Error::new(lit.span(), err).to_compile_error().into(),
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use std::fs;
+
+ use rstest::rstest;
+ use tempfile::TempDir;
+
+ use super::{Error, resolve};
+
+ fn migrations_dir(files: &[&str]) -> TempDir {
+ let dir = TempDir::new().expect("tempdir");
+ for file in files {
+ fs::write(dir.path().join(file), "SELECT 1").expect("write fixture");
+ }
+ dir
+ }
+
+ #[rstest]
+ fn orders_versions_numerically() {
+ let dir = migrations_dir(&["10_tenth.sql", "2_second.sql", "1_first.sql"]);
+ let entries = resolve(dir.path()).expect("resolves");
+ let versions: Vec = entries.iter().map(|entry| entry.version).collect();
+ let descriptions: Vec<&str> = entries
+ .iter()
+ .map(|entry| entry.description.as_str())
+ .collect();
+ assert_eq!(versions, [1, 2, 10]);
+ assert_eq!(descriptions, ["first", "second", "tenth"]);
+ }
+
+ #[rstest]
+ #[case::dash_in_version(&["0001-dash.sql"])]
+ #[case::not_sql(&["notes.txt"])]
+ #[case::empty_description(&["0001_.sql"])]
+ #[case::non_digit_version(&["x_name.sql"])]
+ #[case::uppercase_description(&["0001_Upper.sql"])]
+ #[case::no_underscore(&["0001.sql"])]
+ #[case::plus_sign_version(&["+10_add.sql"])]
+ fn rejects_invalid_names(#[case] files: &[&str]) {
+ let dir = migrations_dir(files);
+ assert!(matches!(
+ resolve(dir.path()),
+ Err(Error::InvalidName { .. })
+ ));
+ }
+
+ #[rstest]
+ fn rejects_subdirectories() {
+ let dir = migrations_dir(&["0001_a.sql"]);
+ fs::create_dir(dir.path().join("0002_b.sql")).expect("subdir");
+ assert!(matches!(
+ resolve(dir.path()),
+ Err(Error::InvalidName { .. })
+ ));
+ }
+
+ #[cfg(unix)]
+ #[rstest]
+ fn rejects_symlinks() {
+ let dir = migrations_dir(&["0001_a.sql"]);
+ let target = TempDir::new().expect("tempdir");
+ let target_file = target.path().join("real.sql");
+ fs::write(&target_file, "SELECT 2").expect("write fixture");
+ std::os::unix::fs::symlink(&target_file, dir.path().join("0002_b.sql")).expect("symlink");
+ assert!(matches!(
+ resolve(dir.path()),
+ Err(Error::InvalidName { .. })
+ ));
+ }
+
+ #[rstest]
+ fn rejects_duplicate_versions() {
+ let dir = migrations_dir(&["0001_a.sql", "1_b.sql"]);
+ assert!(matches!(
+ resolve(dir.path()),
+ Err(Error::DuplicateVersion { version: 1 })
+ ));
+ }
+
+ #[rstest]
+ fn rejects_empty_directory() {
+ let dir = migrations_dir(&[]);
+ assert!(matches!(resolve(dir.path()), Err(Error::Empty { .. })));
+ }
+}
diff --git a/litellm-rust/crates/migrate/Cargo.toml b/litellm-rust/crates/migrate/Cargo.toml
new file mode 100644
index 00000000000..bb1ecaa3128
--- /dev/null
+++ b/litellm-rust/crates/migrate/Cargo.toml
@@ -0,0 +1,12 @@
+[package]
+name = "litellm-migrate"
+version = "0.1.0"
+edition.workspace = true
+license.workspace = true
+repository.workspace = true
+
+[dependencies]
+litellm-migrate-macros.workspace = true
+
+[dev-dependencies]
+rstest.workspace = true
diff --git a/litellm-rust/crates/migrate/README.md b/litellm-rust/crates/migrate/README.md
new file mode 100644
index 00000000000..4817029451c
--- /dev/null
+++ b/litellm-rust/crates/migrate/README.md
@@ -0,0 +1,5 @@
+# Migrations
+
+`litellm-migrate` exports the `Migration` struct and the `migrate!` macro that embeds a directory of `_.sql` files at compile time, sorted by numeric version
+
+The crate does not apply or track migrations; callers decide how and when the embedded SQL runs
diff --git a/litellm-rust/crates/migrate/src/lib.rs b/litellm-rust/crates/migrate/src/lib.rs
new file mode 100644
index 00000000000..f4e065e1b53
--- /dev/null
+++ b/litellm-rust/crates/migrate/src/lib.rs
@@ -0,0 +1,8 @@
+pub use litellm_migrate_macros::migrate;
+
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub struct Migration {
+ pub version: u64,
+ pub description: &'static str,
+ pub sql: &'static str,
+}
diff --git a/litellm-rust/crates/migrate/tests/fixtures/migrations/10_tenth.sql b/litellm-rust/crates/migrate/tests/fixtures/migrations/10_tenth.sql
new file mode 100644
index 00000000000..31807719e9c
--- /dev/null
+++ b/litellm-rust/crates/migrate/tests/fixtures/migrations/10_tenth.sql
@@ -0,0 +1 @@
+SELECT 10;
diff --git a/litellm-rust/crates/migrate/tests/fixtures/migrations/1_first.sql b/litellm-rust/crates/migrate/tests/fixtures/migrations/1_first.sql
new file mode 100644
index 00000000000..e0ac49d1ecf
--- /dev/null
+++ b/litellm-rust/crates/migrate/tests/fixtures/migrations/1_first.sql
@@ -0,0 +1 @@
+SELECT 1;
diff --git a/litellm-rust/crates/migrate/tests/fixtures/migrations/2_second.sql b/litellm-rust/crates/migrate/tests/fixtures/migrations/2_second.sql
new file mode 100644
index 00000000000..e7f8100648d
--- /dev/null
+++ b/litellm-rust/crates/migrate/tests/fixtures/migrations/2_second.sql
@@ -0,0 +1 @@
+SELECT 2;
diff --git a/litellm-rust/crates/migrate/tests/migrate.rs b/litellm-rust/crates/migrate/tests/migrate.rs
new file mode 100644
index 00000000000..61c80351cf4
--- /dev/null
+++ b/litellm-rust/crates/migrate/tests/migrate.rs
@@ -0,0 +1,21 @@
+use litellm_migrate::Migration;
+use rstest::rstest;
+
+const MIGRATIONS: &[Migration] = litellm_migrate::migrate!("tests/fixtures/migrations");
+
+#[rstest]
+#[case::first(0, 1, "first", include_str!("fixtures/migrations/1_first.sql"))]
+#[case::second(1, 2, "second", include_str!("fixtures/migrations/2_second.sql"))]
+#[case::tenth(2, 10, "tenth", include_str!("fixtures/migrations/10_tenth.sql"))]
+fn embeds_every_file_sorted_by_numeric_version(
+ #[case] index: usize,
+ #[case] version: u64,
+ #[case] description: &str,
+ #[case] sql: &str,
+) {
+ assert_eq!(MIGRATIONS.len(), 3);
+ let migration = &MIGRATIONS[index];
+ assert_eq!(migration.version, version);
+ assert_eq!(migration.description, description);
+ assert_eq!(migration.sql, sql);
+}
diff --git a/litellm-rust/crates/python-bridge/src/lib.rs b/litellm-rust/crates/python-bridge/src/lib.rs
index d269fa4015f..6e091f0594f 100644
--- a/litellm-rust/crates/python-bridge/src/lib.rs
+++ b/litellm-rust/crates/python-bridge/src/lib.rs
@@ -44,7 +44,10 @@ mod _native {
#[pymodule_export]
use crate::routes::token_counter::TokenCounter;
#[pymodule_export]
- use crate::routes::traces::{NativeTraceStorage, trace_decode_otlp, trace_encode_error};
+ use crate::routes::traces::{
+ NativeTraceStorage, trace_decode_otlp, trace_encode_error,
+ trace_normalized_field_definitions,
+ };
#[cfg(feature = "huggingface")]
#[pymodule_export]
use crate::tokenizer::HuggingFaceEncoding;
@@ -112,6 +115,7 @@ mod tests {
"NativeTraceStorage",
"trace_decode_otlp",
"trace_encode_error",
+ "trace_normalized_field_definitions",
"TokenCounter",
"Tokenizer",
"gil_stats",
diff --git a/litellm-rust/crates/python-bridge/src/routes/traces.rs b/litellm-rust/crates/python-bridge/src/routes/traces.rs
index ca66e2e46be..0ff5fe98f55 100644
--- a/litellm-rust/crates/python-bridge/src/routes/traces.rs
+++ b/litellm-rust/crates/python-bridge/src/routes/traces.rs
@@ -3,7 +3,9 @@ use std::collections::BTreeMap;
use litellm_host_python::{FromPythonCache, ToPythonCache};
use litellm_http::ClientVariant;
use litellm_storage_clickhouse::Storage;
-use litellm_traces::{Error, InsertTable, Parameter, ReadQuery, Shared};
+use litellm_traces::{
+ Error, InsertTable, Parameter, QueryAccessError, QueryReaders, QueryScope, ReadQuery, Shared,
+};
use prost::Message;
use pyo3::{
exceptions::{PyOverflowError, PyRuntimeError, PyValueError},
@@ -46,9 +48,25 @@ fn map_error(error: Error) -> PyErr {
}
}
+fn map_sql_error(error: Error) -> PyErr {
+ match error {
+ Error::QueryFailed(400 | 404) => PyValueError::new_err(error.to_string()),
+ error => map_error(error),
+ }
+}
+
+fn map_query_access_error(error: QueryAccessError) -> PyErr {
+ match error {
+ QueryAccessError::Storage(error) => map_sql_error(error),
+ QueryAccessError::InvalidScope => PyValueError::new_err(error.to_string()),
+ error => PyRuntimeError::new_err(error.to_string()),
+ }
+}
+
#[pyclass]
pub struct NativeTraceStorage {
storage: Storage,
+ query_readers: QueryReaders,
}
#[pymethods]
@@ -57,8 +75,13 @@ impl NativeTraceStorage {
#[pyo3(signature = (database, url, reader_url = None))]
fn new(database: String, url: &str, reader_url: Option<&str>) -> PyResult {
litellm_traces::schema_statements(&database, 1, 1).map_err(map_error)?;
+ let storage = Storage::new(database, url, reader_url).map_err(map_error)?;
Ok(Self {
- storage: Storage::new(database, url, reader_url).map_err(map_error)?,
+ query_readers: QueryReaders::new(
+ storage.writer().clone(),
+ storage.database().to_owned(),
+ ),
+ storage,
})
}
@@ -107,6 +130,52 @@ impl NativeTraceStorage {
)
}
+ fn query_sql<'py>(
+ &self,
+ py: Python<'py>,
+ sql: String,
+ #[pyo3(from_py_with = litellm_host_python::from_py_argument)] scope: QueryScope,
+ secret: String,
+ ) -> PyResult> {
+ if sql.trim().is_empty() {
+ return Err(map_error(Error::EmptySql));
+ }
+ let readers = self.query_readers.clone();
+ let client = crate::http::host_client(py, ClientVariant::NoRedirect)?;
+ crate::execution::run_async(
+ py,
+ async move {
+ let _permit = readers.acquire()?;
+ let connection = readers.connection(&client, &scope, &secret).await?;
+ litellm_traces::query_sql(&client, &connection, &sql)
+ .await
+ .map_err(QueryAccessError::Storage)
+ },
+ map_query_access_error,
+ )
+ }
+
+ fn query_help<'py>(
+ &self,
+ py: Python<'py>,
+ #[pyo3(from_py_with = litellm_host_python::from_py_argument)] scope: QueryScope,
+ secret: String,
+ ) -> PyResult> {
+ let readers = self.query_readers.clone();
+ let client = crate::http::host_client(py, ClientVariant::NoRedirect)?;
+ crate::execution::run_async(
+ py,
+ async move {
+ let _permit = readers.acquire()?;
+ let connection = readers.connection(&client, &scope, &secret).await?;
+ litellm_traces::query_help(&client, &connection)
+ .await
+ .map_err(QueryAccessError::Storage)
+ },
+ map_query_access_error,
+ )
+ }
+
fn lens_query<'py>(
&self,
py: Python<'py>,
@@ -236,6 +305,14 @@ fn spans_to_py<'py>(
"events",
litellm_host_python::Pythonized(&span.events).into_pyobject(py)?,
)?;
+ row.set_item(
+ "normalized",
+ litellm_host_python::Pythonized(&span.normalized).into_pyobject(py)?,
+ )?;
+ row.set_item(
+ "consumed_attributes",
+ litellm_host_python::Pythonized(&span.consumed_attributes).into_pyobject(py)?,
+ )?;
result.append(row)?;
}
Ok(result)
@@ -288,3 +365,8 @@ mod tests {
});
}
}
+
+#[pyfunction]
+pub fn trace_normalized_field_definitions<'py>(py: Python<'py>) -> PyResult> {
+ litellm_host_python::Pythonized(litellm_traces::NORMALIZED_FIELD_DEFINITIONS).into_pyobject(py)
+}
diff --git a/litellm-rust/crates/traces/AGENTS.md b/litellm-rust/crates/traces/AGENTS.md
index 645e88dfae1..d0181c53308 100644
--- a/litellm-rust/crates/traces/AGENTS.md
+++ b/litellm-rust/crates/traces/AGENTS.md
@@ -1,6 +1,6 @@
- Keep OTLP decoding, trace schema, row encoding and named query selection here. Generic ClickHouse connections and HTTP execution belong in `litellm-storage-clickhouse`
- Keep this crate independent of Python; PyO3 conversion and public Python exceptions belong in `python-bridge`
-- Keep the SQL migrations here as the only ClickHouse schema definition
+- Keep the SQL migrations here as the only ClickHouse schema definition, as `migrations/NNNN_description.sql` files embedded by `litellm_migrate::migrate!`; adding a file is the only step
- Use typed query parameters and a dedicated SELECT-only reader with server-side limits
- Keep `config/reader.xml` grants on the database the schema is created in (CLICKHOUSE_DATABASE, default `litellm`)
- Bound insert time and encoded bytes; make retry deduplication behavior explicit for supported ClickHouse versions
diff --git a/litellm-rust/crates/traces/Cargo.toml b/litellm-rust/crates/traces/Cargo.toml
index 74de400764c..5e3b41e719a 100644
--- a/litellm-rust/crates/traces/Cargo.toml
+++ b/litellm-rust/crates/traces/Cargo.toml
@@ -6,18 +6,26 @@ license.workspace = true
repository.workspace = true
[dependencies]
+askama.workspace = true
base64.workspace = true
flate2.workspace = true
+futures-util.workspace = true
+hmac = "0.12.1"
+indexmap = { version = "2", features = ["serde"] }
+moka.workspace = true
opentelemetry-proto = { workspace = true, features = ["gen-tonic-messages", "trace", "with-serde"] }
prost.workspace = true
time = { workspace = true, features = ["formatting"] }
litellm-http.workspace = true
+litellm-migrate.workspace = true
litellm-storage-clickhouse.workspace = true
sha2.workspace = true
serde = { workspace = true, features = ["rc"] }
serde_json.workspace = true
strum.workspace = true
thiserror.workspace = true
+tokio.workspace = true
+url.workspace = true
[dev-dependencies]
criterion.workspace = true
diff --git a/litellm-rust/crates/traces/build.rs b/litellm-rust/crates/traces/build.rs
new file mode 100644
index 00000000000..3a8149ef075
--- /dev/null
+++ b/litellm-rust/crates/traces/build.rs
@@ -0,0 +1,3 @@
+fn main() {
+ println!("cargo:rerun-if-changed=migrations");
+}
diff --git a/litellm-rust/crates/traces/migrations/0001_otel_traces.sql b/litellm-rust/crates/traces/migrations/0001_otel_traces.sql
index d8e0184b5a3..fb5eaa367d7 100644
--- a/litellm-rust/crates/traces/migrations/0001_otel_traces.sql
+++ b/litellm-rust/crates/traces/migrations/0001_otel_traces.sql
@@ -38,10 +38,11 @@ CREATE TABLE IF NOT EXISTS {database}.otel_traces
Input String CODEC(ZSTD(3)),
Output String CODEC(ZSTD(3)),
InputPreview String DEFAULT substring(Input, 1, 240),
+ EngineReceivedMs UInt64 DEFAULT 0,
INDEX idx_trace_id TraceId TYPE bloom_filter(0.001) GRANULARITY 1,
INDEX idx_req_id LiteLLMRequestId TYPE bloom_filter(0.01) GRANULARITY 1
)
ENGINE = MergeTree
PARTITION BY toDate(Timestamp)
ORDER BY (TeamId, ServiceName, toDateTime(Timestamp), TraceId)
-SETTINGS ttl_only_drop_parts = 1, non_replicated_deduplication_window = 1000
+SETTINGS ttl_only_drop_parts = 1, materialize_ttl_recalculate_only = 1, non_replicated_deduplication_window = 1000
diff --git a/litellm-rust/crates/traces/migrations/0005_otel_traces_ttl.sql b/litellm-rust/crates/traces/migrations/0002_otel_traces_ttl.sql
similarity index 100%
rename from litellm-rust/crates/traces/migrations/0005_otel_traces_ttl.sql
rename to litellm-rust/crates/traces/migrations/0002_otel_traces_ttl.sql
diff --git a/litellm-rust/crates/traces/migrations/0002_agent_traces.sql b/litellm-rust/crates/traces/migrations/0003_agent_traces.sql
similarity index 93%
rename from litellm-rust/crates/traces/migrations/0002_agent_traces.sql
rename to litellm-rust/crates/traces/migrations/0003_agent_traces.sql
index 0c3547872bb..821cc2f3723 100644
--- a/litellm-rust/crates/traces/migrations/0002_agent_traces.sql
+++ b/litellm-rust/crates/traces/migrations/0003_agent_traces.sql
@@ -22,4 +22,4 @@ CREATE TABLE IF NOT EXISTS {database}.agent_traces_by_key
)
ENGINE = AggregatingMergeTree
ORDER BY (TeamId, ApiKeyHash, TraceId)
-SETTINGS non_replicated_deduplication_window = 1000
+SETTINGS materialize_ttl_recalculate_only = 1, non_replicated_deduplication_window = 1000
diff --git a/litellm-rust/crates/traces/migrations/0006_agent_traces_ttl.sql b/litellm-rust/crates/traces/migrations/0004_agent_traces_ttl.sql
similarity index 100%
rename from litellm-rust/crates/traces/migrations/0006_agent_traces_ttl.sql
rename to litellm-rust/crates/traces/migrations/0004_agent_traces_ttl.sql
diff --git a/litellm-rust/crates/traces/migrations/0003_agent_traces_mv.sql b/litellm-rust/crates/traces/migrations/0005_agent_traces_mv.sql
similarity index 100%
rename from litellm-rust/crates/traces/migrations/0003_agent_traces_mv.sql
rename to litellm-rust/crates/traces/migrations/0005_agent_traces_mv.sql
diff --git a/litellm-rust/crates/traces/migrations/0004_spend_logs.sql b/litellm-rust/crates/traces/migrations/0006_spend_logs.sql
similarity index 94%
rename from litellm-rust/crates/traces/migrations/0004_spend_logs.sql
rename to litellm-rust/crates/traces/migrations/0006_spend_logs.sql
index a14930f438f..44f7959b2bf 100644
--- a/litellm-rust/crates/traces/migrations/0004_spend_logs.sql
+++ b/litellm-rust/crates/traces/migrations/0006_spend_logs.sql
@@ -34,9 +34,11 @@ CREATE TABLE IF NOT EXISTS {database}.spend_logs
metadata String CODEC(ZSTD(3)),
messages String CODEC(ZSTD(3)),
response String CODEC(ZSTD(3)),
+ EngineReceivedMs UInt64 DEFAULT 0,
INDEX idx_response_id response_id TYPE bloom_filter(0.001) GRANULARITY 1,
INDEX idx_trace_id trace_id TYPE bloom_filter(0.001) GRANULARITY 1
)
ENGINE = ReplacingMergeTree(end_time)
PARTITION BY toYYYYMM(start_time)
ORDER BY (team_id, start_time, request_id)
+SETTINGS materialize_ttl_recalculate_only = 1
diff --git a/litellm-rust/crates/traces/migrations/0008_trace_received.sql b/litellm-rust/crates/traces/migrations/0008_trace_received.sql
deleted file mode 100644
index 9d8113b2430..00000000000
--- a/litellm-rust/crates/traces/migrations/0008_trace_received.sql
+++ /dev/null
@@ -1 +0,0 @@
-ALTER TABLE {database}.otel_traces ADD COLUMN IF NOT EXISTS EngineReceivedMs UInt64 DEFAULT 0
diff --git a/litellm-rust/crates/traces/migrations/0009_spend_received.sql b/litellm-rust/crates/traces/migrations/0009_spend_received.sql
deleted file mode 100644
index 2b2d2c7e5d7..00000000000
--- a/litellm-rust/crates/traces/migrations/0009_spend_received.sql
+++ /dev/null
@@ -1 +0,0 @@
-ALTER TABLE {database}.spend_logs ADD COLUMN IF NOT EXISTS EngineReceivedMs UInt64 DEFAULT 0
diff --git a/litellm-rust/crates/traces/query/lens_agents.sql b/litellm-rust/crates/traces/query/lens_agents.sql
new file mode 100644
index 00000000000..fbdd578f8e7
--- /dev/null
+++ b/litellm-rust/crates/traces/query/lens_agents.sql
@@ -0,0 +1,6 @@
+SELECT DISTINCT AgentName AS agent_name
+FROM otel_traces
+WHERE AgentName != ''
+ AND ({all_teams:UInt8}=1 OR TeamId={team:String})
+ AND ({key_hash:String}='' OR ApiKeyHash={key_hash:String})
+ORDER BY agent_name
diff --git a/litellm-rust/crates/traces/query/lens_availability.sql b/litellm-rust/crates/traces/query/lens_availability.sql
new file mode 100644
index 00000000000..8d350dd1779
--- /dev/null
+++ b/litellm-rust/crates/traces/query/lens_availability.sql
@@ -0,0 +1,8 @@
+SELECT
+ EXISTS(SELECT 1 FROM otel_traces
+ WHERE ({all_teams:UInt8}=1 OR TeamId={team:String})
+ AND ({key_hash:String}='' OR ApiKeyHash={key_hash:String})) AS traces,
+ EXISTS(SELECT 1 FROM spend_logs
+ WHERE ({all_teams:UInt8}=1 OR team_id={team:String})
+ AND ({key_hash:String}='' OR api_key={key_hash:String})
+ AND NOT JSONExtractBool(metadata,'litellm_lens_internal')) AS requests
diff --git a/litellm-rust/crates/traces/query/lens_sample.sql b/litellm-rust/crates/traces/query/lens_sample.sql
index 1fc9c964a6f..92086c33c13 100644
--- a/litellm-rust/crates/traces/query/lens_sample.sql
+++ b/litellm-rust/crates/traces/query/lens_sample.sql
@@ -28,6 +28,7 @@ SELECT *, selection_key FROM (
GROUP BY TeamId,ApiKeyHash,TraceId
HAVING max(EngineReceivedMs) < {end:UInt64}
AND max(toUnixTimestamp64Milli(Timestamp)+toInt64(intDiv(Duration,1000000))) < {end:UInt64}
+ AND ({agent_name:String}='' OR countIf(AgentName={agent_name:String}) > 0)
AND countIf(arrayAll((k,v) -> ResourceAttributes[k]=v OR SpanAttributes[k]=v,
{filter_keys:Array(String)},{filter_values:Array(String)})
AND ({service:String}='' OR ServiceName={service:String})) > 0
@@ -48,6 +49,7 @@ SELECT *, selection_key FROM (
OR JSONExtractString(metadata,'requester_metadata',k)=v OR (k='tag' AND has(request_tags,v)),
{filter_keys:Array(String)},{filter_values:Array(String)})
AND ({service:String}='' OR model_group={service:String})
+ AND {agent_name:String}=''
AND NOT JSONExtractBool(metadata,'litellm_lens_internal')
AND ({source:String}!='both' OR (team_id,api_key,response_id) NOT IN (
SELECT TeamId,ApiKeyHash,LiteLLMRequestId FROM otel_traces
diff --git a/litellm-rust/crates/traces/src/error.rs b/litellm-rust/crates/traces/src/error.rs
index 18fa4af9b53..c677e73e96f 100644
--- a/litellm-rust/crates/traces/src/error.rs
+++ b/litellm-rust/crates/traces/src/error.rs
@@ -4,4 +4,26 @@ pub enum DecodeError {
InvalidPayload,
#[error("OTLP trace payload exceeds the decoding budget")]
TooLarge,
+ #[error("OTLP token count is outside the storage range")]
+ TokenCountOutOfRange,
+}
+
+#[derive(Debug, thiserror::Error)]
+pub enum QueryAccessError {
+ #[error("trace SQL queries require a configured proxy master key")]
+ MissingSecret,
+ #[error("invalid trace query scope")]
+ InvalidScope,
+ #[error("trace SQL query concurrency limit exceeded")]
+ Busy,
+ #[error(
+ "ClickHouse reader provisioning failed with HTTP status {0}; the configured connection must be allowed to manage users, row policies, and SELECT grants on the trace tables"
+ )]
+ ProvisionFailed(u16),
+ #[error("ClickHouse reader provisioning transport failed")]
+ ProvisionTransport,
+ #[error(transparent)]
+ Storage(#[from] litellm_storage_clickhouse::Error),
+ #[error(transparent)]
+ Cached(#[from] std::sync::Arc),
}
diff --git a/litellm-rust/crates/traces/src/lib.rs b/litellm-rust/crates/traces/src/lib.rs
index 1489b44c118..9e01d10f73e 100644
--- a/litellm-rust/crates/traces/src/lib.rs
+++ b/litellm-rust/crates/traces/src/lib.rs
@@ -1,14 +1,23 @@
mod error;
mod insert;
+mod normalize;
mod otlp;
+mod query;
+mod query_access;
mod schema;
mod shared;
mod sql;
-pub use error::DecodeError;
+pub use error::{DecodeError, QueryAccessError};
pub use insert::{InsertRow, InsertTable, encode_rows, insert_rows, insert_shared_rows};
pub use litellm_storage_clickhouse::{Connection, Error, Parameter, execute_read};
+pub use normalize::{
+ NORMALIZED_FIELD_DEFINITIONS, NormalizedFieldDefinition, NormalizedSpan, ObservationType,
+};
pub use otlp::{DecodedSpan, decode_otlp};
+pub use query_access::{QueryReaders, QueryScope};
pub use schema::{ensure_schema, schema_statements};
pub use shared::{Shared, SharedIdentity};
pub use sql::{LensQuery, ReadQuery, execute_named_read};
+
+pub use query::{query_help, query_sql};
diff --git a/litellm-rust/crates/traces/src/normalize/genai.rs b/litellm-rust/crates/traces/src/normalize/genai.rs
new file mode 100644
index 00000000000..15a334ab0fd
--- /dev/null
+++ b/litellm-rust/crates/traces/src/normalize/genai.rs
@@ -0,0 +1,63 @@
+use std::collections::BTreeMap;
+
+use super::{NormalizedSpan, ObservationType, SpanNormalizer, attr, first, usage_tokens};
+use crate::DecodeError;
+
+pub(super) struct GenAiNormalizer;
+
+impl SpanNormalizer for GenAiNormalizer {
+ fn matches(&self, _scope_name: &str, _attributes: &BTreeMap) -> bool {
+ true
+ }
+
+ fn consumed_attributes(&self, attributes: &BTreeMap) -> [&'static str; 2] {
+ [
+ if attr(attributes, "gen_ai.input.messages").is_empty() {
+ "gen_ai.tool.call.arguments"
+ } else {
+ "gen_ai.input.messages"
+ },
+ if attr(attributes, "gen_ai.output.messages").is_empty() {
+ "gen_ai.tool.call.result"
+ } else {
+ "gen_ai.output.messages"
+ },
+ ]
+ }
+
+ fn normalize(
+ &self,
+ _name: &str,
+ parent_span_id: &str,
+ attributes: &BTreeMap,
+ ) -> Result {
+ let (input_tokens, output_tokens) = usage_tokens(attributes)?;
+ let observation_type = match attr(attributes, "gen_ai.operation.name") {
+ "invoke_agent" => ObservationType::Agent,
+ "chat" | "text_completion" | "generate_content" => ObservationType::Llm,
+ "execute_tool" => ObservationType::Tool,
+ _ if parent_span_id.is_empty() => ObservationType::Agent,
+ _ => ObservationType::Chain,
+ };
+ Ok(NormalizedSpan {
+ observation_type,
+ agent_name: attr(attributes, "gen_ai.agent.name").to_owned(),
+ litellm_request_id: attr(attributes, "gen_ai.response.id").to_owned(),
+ model: first(attributes, "gen_ai.request.model", "gen_ai.response.model").to_owned(),
+ input_tokens,
+ output_tokens,
+ input: first(
+ attributes,
+ "gen_ai.input.messages",
+ "gen_ai.tool.call.arguments",
+ )
+ .to_owned(),
+ output: first(
+ attributes,
+ "gen_ai.output.messages",
+ "gen_ai.tool.call.result",
+ )
+ .to_owned(),
+ })
+ }
+}
diff --git a/litellm-rust/crates/traces/src/normalize/langsmith.rs b/litellm-rust/crates/traces/src/normalize/langsmith.rs
new file mode 100644
index 00000000000..bfe7a2c1796
--- /dev/null
+++ b/litellm-rust/crates/traces/src/normalize/langsmith.rs
@@ -0,0 +1,468 @@
+use std::{collections::BTreeMap, io};
+
+use indexmap::IndexMap;
+use serde::{Deserialize, Deserializer, Serialize, de::DeserializeOwned};
+use serde_json::{Value, ser::Formatter};
+
+use super::{NormalizedSpan, ObservationType, SpanNormalizer, attr, usage_tokens};
+use crate::DecodeError;
+
+pub(super) struct LangSmithNormalizer;
+
+#[derive(Deserialize)]
+#[serde(untagged)]
+enum MessageContent {
+ Text(String),
+ Blocks(Vec),
+ Other(Value),
+}
+
+impl MessageContent {
+ fn display_text(&self) -> String {
+ match self {
+ Self::Text(text) => text.clone(),
+ Self::Blocks(blocks) => blocks
+ .iter()
+ .filter_map(|block| match block {
+ ContentBlock::Text { text } => Some(text.as_str()),
+ ContentBlock::Hidden(kind) => match kind {
+ HiddenBlock::Reasoning
+ | HiddenBlock::Thinking
+ | HiddenBlock::RedactedThinking
+ | HiddenBlock::FunctionCall
+ | HiddenBlock::ToolUse
+ | HiddenBlock::ToolCall => None,
+ },
+ })
+ .collect::>()
+ .join("\n\n"),
+ Self::Other(value) => encode(value),
+ }
+ }
+}
+
+#[derive(Deserialize)]
+#[serde(untagged)]
+enum ContentBlock {
+ Text { text: String },
+ Hidden(HiddenBlock),
+}
+
+#[derive(Deserialize)]
+#[serde(tag = "type", rename_all = "snake_case")]
+enum HiddenBlock {
+ Reasoning,
+ Thinking,
+ RedactedThinking,
+ FunctionCall,
+ ToolUse,
+ ToolCall,
+}
+
+#[derive(Deserialize, Serialize)]
+#[serde(transparent)]
+struct RawToolCall(IndexMap);
+
+#[derive(Deserialize)]
+struct ResponseMetadata {
+ id: Option,
+}
+
+#[derive(Deserialize)]
+struct RawMessage {
+ kwargs: Option>,
+ #[serde(rename = "type")]
+ kind: Option,
+ role: Option,
+ content: Option,
+ tool_calls: Option>,
+ name: Option,
+ response_metadata: Option,
+}
+
+impl RawMessage {
+ fn unwrapped(&self) -> &Self {
+ self.kwargs.as_deref().unwrap_or(self)
+ }
+
+ fn normalized(&self) -> NormalizedMessage<'_> {
+ let fields = self.unwrapped();
+ let raw_role = fields
+ .kind
+ .as_deref()
+ .filter(|role| !role.is_empty())
+ .or_else(|| fields.role.as_deref().filter(|role| !role.is_empty()))
+ .unwrap_or_default();
+ let role = match raw_role {
+ "human" => "user",
+ "ai" => "assistant",
+ other => other,
+ };
+ NormalizedMessage {
+ role,
+ content: fields
+ .content
+ .as_ref()
+ .map_or_else(String::new, MessageContent::display_text),
+ tool_calls: fields
+ .tool_calls
+ .as_deref()
+ .filter(|calls| !calls.is_empty()),
+ name: (role == "tool")
+ .then_some(fields.name.as_ref())
+ .flatten()
+ .filter(|name| !name.is_null() && name != &&Value::String(String::new())),
+ }
+ }
+}
+
+#[derive(Serialize)]
+struct NormalizedMessage<'a> {
+ role: &'a str,
+ content: String,
+ #[serde(skip_serializing_if = "Option::is_none")]
+ tool_calls: Option<&'a [RawToolCall]>,
+ #[serde(skip_serializing_if = "Option::is_none")]
+ name: Option<&'a Value>,
+}
+
+enum MessageBatch {
+ Flat(Vec),
+ Nested(Vec>),
+}
+
+impl<'de> Deserialize<'de> for MessageBatch {
+ fn deserialize>(deserializer: D) -> Result {
+ let value = Value::deserialize(deserializer)?;
+ let Value::Array(items) = value else {
+ return Err(serde::de::Error::custom("messages must be an array"));
+ };
+ let parse = |items: Vec| {
+ items
+ .into_iter()
+ .filter_map(|item| serde_json::from_value(item).ok())
+ .collect()
+ };
+ Ok(if items.first().is_some_and(Value::is_array) {
+ Self::Nested(
+ items
+ .into_iter()
+ .filter_map(|item| item.as_array().cloned())
+ .map(parse)
+ .collect(),
+ )
+ } else {
+ Self::Flat(parse(items))
+ })
+ }
+}
+
+fn lenient<'de, D: Deserializer<'de>, T: DeserializeOwned>(
+ deserializer: D,
+) -> Result