import glob import os import random import re import shutil import subprocess import tempfile import time from dataclasses import dataclass, replace from pathlib import Path from typing import Optional from litellm_proxy_extras._logging import logger from litellm_proxy_extras.replica_identity import ( REPLICA_IDENTITY_FULL_ENV_VAR, apply_replica_identity_full, ) from litellm_proxy_extras.prisma_toolchain import ( PRISMA_COMMAND_TIMEOUT_ENV_VAR, PRISMA_MIGRATE_DEPLOY_TIMEOUT_ENV_VAR, ensure_prisma_toolchain, prisma_command_timeout, prisma_migrate_deploy_timeout, ) def str_to_bool(value: Optional[str]) -> bool: if value is None: return False return value.lower() in ("true", "1", "t", "y", "yes") def _get_prisma_env() -> dict: """Get environment variables for Prisma, handling offline mode if configured.""" prisma_env = os.environ.copy() if str_to_bool(os.getenv("PRISMA_OFFLINE_MODE")): # These env vars prevent Prisma from attempting downloads prisma_env["NPM_CONFIG_PREFER_OFFLINE"] = "true" prisma_env["NPM_CONFIG_CACHE"] = os.getenv( "NPM_CONFIG_CACHE", "/app/.cache/npm" ) return prisma_env _MIGRATION_TS_RE = re.compile(r"^(\d{14})_") _MIGRATION_DEADLOCK_MARKER = "deadlock detected" MAX_MIGRATE_DEPLOY_ATTEMPTS = 4 @dataclass(frozen=True) class _MigrateAttemptBudget: """Retries left, and the recoveries already run. A recovery that lands something new costs nothing, so a database full of objects `prisma db push` created works through them one per pass. Anything that made no progress spends an attempt, so a stuck run still gives up. """ attempts_left: int recoveries: frozenset[str] = frozenset() @property def exhausted(self) -> bool: return self.attempts_left <= 0 @property def attempt_number(self) -> int: return MAX_MIGRATE_DEPLOY_ATTEMPTS - self.attempts_left + 1 def spend(self) -> "_MigrateAttemptBudget": return replace(self, attempts_left=self.attempts_left - 1) def after_recovery(self, recovery: str) -> "_MigrateAttemptBudget": if recovery in self.recoveries: return self.spend() return replace(self, recoveries=self.recoveries | {recovery}) _SPEND_LOGS_ALTER_RE = re.compile(r'^ALTER\s+TABLE\s+"LiteLLM_SpendLogs"\s', re.IGNORECASE) _SPEND_LOGS_ARTIFACT_DROP_RE = re.compile( r'^DROP\s+TABLE\s+"LiteLLM_SpendLogs_[^"]*"', re.IGNORECASE ) _SPEND_LOGS_PK_CLAUSE_RE = re.compile( r'^(?:DROP\s+CONSTRAINT\s+"[^"]*_pkey"' r'|ADD\s+(?:CONSTRAINT\s+"[^"]*"\s+)?PRIMARY\s+KEY\s*\([^)]*\))$', re.IGNORECASE, ) PARTITIONED_SPEND_LOGS_PUSH_ERROR = ( "LiteLLM_SpendLogs is a partitioned table (see db_scripts/partition_spend_logs.sql), " "so its primary key must include the partition key (\"startTime\"). `prisma db push` " "reconciles the database against schema.prisma, which declares the unpartitioned " "primary key (\"request_id\"), and Postgres rejects that rewrite with: unique " "constraint on partitioned table must include all partitioning columns. Start the " "proxy without --use_prisma_db_push so it uses `prisma migrate deploy`, which only " "applies shipped migrations and leaves the partitioned primary key alone." ) def _without_sql_comments(statement: str) -> str: return "\n".join( line for line in statement.splitlines() if line.strip() and not line.strip().startswith("--") ).strip() def _without_spend_logs_pk_clauses(statement: str) -> Optional[str]: prefix_match = _SPEND_LOGS_ALTER_RE.match(statement) if not prefix_match: return statement kept = tuple( clause.strip() for clause in statement[prefix_match.end():].split(",\n") if not _SPEND_LOGS_PK_CLAUSE_RE.match(clause.strip()) ) if not kept: return None return statement[: prefix_match.end()] + ",\n".join(kept) def filter_partitioned_spend_logs_diff(diff_sql: str) -> str: """Drop statements from a `prisma migrate diff` script that fight the SpendLogs partitioning runbook (db_scripts/partition_spend_logs.sql): the primary-key rewrite on "LiteLLM_SpendLogs", which Postgres rejects on a partitioned table, and drops of runbook artifacts such as "LiteLLM_SpendLogs_legacy".""" kept = tuple( filtered for statement in diff_sql.split(";") for bare in (_without_sql_comments(statement),) if bare and not _SPEND_LOGS_ARTIFACT_DROP_RE.match(bare) for filtered in (_without_spend_logs_pk_clauses(bare),) if filtered is not None ) return "".join(f"{statement};\n\n" for statement in kept) def _migration_timestamp(name: str) -> int: """Extract the leading `YYYYMMDDHHMMSS` timestamp from a migration name. Returns 0 if the name doesn't match the Prisma pattern — unexpected-format entries sort as "oldest" and are treated as historical. """ m = _MIGRATION_TS_RE.match(name) return int(m.group(1)) if m else 0 def _max_migration_timestamp(names) -> int: """Max timestamp in a set/list of migration names (0 if empty).""" if not names: return 0 return max(_migration_timestamp(n) for n in names) def _get_prisma_command() -> str: """Get the Prisma command to use, bypassing Python wrapper in offline mode.""" if str_to_bool(os.getenv("PRISMA_OFFLINE_MODE")): # Primary location where Prisma Python package installs the CLI default_cli_path = "/app/.cache/prisma-python/binaries/node_modules/.bin/prisma" # Check if custom path is provided (for flexibility) custom_cli_path = os.getenv("PRISMA_CLI_PATH") if custom_cli_path and os.path.exists(custom_cli_path): logger.info(f"Using custom Prisma CLI at {custom_cli_path}") return custom_cli_path # Check the default location if os.path.exists(default_cli_path): logger.info(f"Using cached Prisma CLI at {default_cli_path}") return default_cli_path # If not found, log warning and fall back logger.warning( f"Prisma CLI not found at {default_cli_path}. " "Falling back to Python wrapper (may attempt downloads)" ) # Fall back to the Python wrapper (will work in online mode) return "prisma" class ProxyExtrasDBManager: @staticmethod def _get_prisma_dir() -> str: """ Get the path to the migrations directory Set os.environ["LITELLM_MIGRATION_DIR"] to a custom migrations directory, to support baselining db in read-only fs. """ custom_migrations_dir = os.getenv("LITELLM_MIGRATION_DIR") pkg_migrations_dir = os.path.dirname(__file__) if custom_migrations_dir: # If migrations_dir exists, copy contents if os.path.exists(custom_migrations_dir): # Copy contents instead of directory itself for item in os.listdir(pkg_migrations_dir): src_path = os.path.join(pkg_migrations_dir, item) dst_path = os.path.join(custom_migrations_dir, item) if os.path.isdir(src_path): shutil.copytree(src_path, dst_path, dirs_exist_ok=True) else: shutil.copy2(src_path, dst_path) else: # If directory doesn't exist, create it and copy everything shutil.copytree(pkg_migrations_dir, custom_migrations_dir) return custom_migrations_dir return pkg_migrations_dir @staticmethod def _create_baseline_migration(schema_path: str) -> bool: """Create a baseline migration for an existing database""" prisma_dir = ProxyExtrasDBManager._get_prisma_dir() prisma_dir_path = Path(prisma_dir) init_dir = prisma_dir_path / "migrations" / "0_init" # Create migrations/0_init directory init_dir.mkdir(parents=True, exist_ok=True) database_url = os.getenv("DATABASE_URL") if not database_url: logger.error("DATABASE_URL not set") return False # Set up environment for offline mode if configured prisma_env = _get_prisma_env() try: # 1. Generate migration SQL file by comparing empty state to current db state logger.info("Generating baseline migration...") migration_file = init_dir / "migration.sql" subprocess.run( [ _get_prisma_command(), "migrate", "diff", "--from-empty", "--to-url", database_url, "--script", ], stdout=open(migration_file, "w"), check=True, timeout=prisma_command_timeout(), env=prisma_env, ) # 3. Mark the migration as applied since it represents current state logger.info("Marking baseline migration as applied...") subprocess.run( [ _get_prisma_command(), "migrate", "resolve", "--applied", "0_init", ], check=True, timeout=prisma_command_timeout(), env=prisma_env, ) return True except subprocess.TimeoutExpired: logger.warning( "Migration timed out - the database might be under heavy load." ) return False except subprocess.CalledProcessError as e: logger.warning( f"Error creating baseline migration: {e}, {e.stderr}, {e.stdout}" ) raise e @staticmethod def _get_migration_names(migrations_dir: str) -> list: """Get all migration directory names from the migrations folder""" migration_paths = glob.glob(f"{migrations_dir}/migrations/*/migration.sql") logger.info(f"Found {len(migration_paths)} migrations at {migrations_dir}") return [Path(p).parent.name for p in migration_paths] @staticmethod def _roll_back_migration(migration_name: str): """Mark a specific migration as rolled back""" # Set up environment for offline mode if configured prisma_env = _get_prisma_env() subprocess.run( [ _get_prisma_command(), "migrate", "resolve", "--rolled-back", migration_name, ], timeout=prisma_command_timeout(), check=True, capture_output=True, env=prisma_env, ) @staticmethod def _roll_back_migration_best_effort(migration_name: str) -> None: """Mark a migration rolled back, tolerating a concurrent resolver having already done it.""" try: ProxyExtrasDBManager._roll_back_migration(migration_name) except (subprocess.CalledProcessError, subprocess.TimeoutExpired): pass @staticmethod def _failed_migration_logs(migration_name: str) -> Optional[str]: """Return failed migration logs, or None if the ledger is unavailable.""" database_url = os.getenv("DATABASE_URL") if not database_url: return None try: import psycopg except ImportError: return None cleaned_url = ProxyExtrasDBManager._strip_prisma_query_params(database_url) ledger_table = psycopg.sql.SQL("{}.{}").format( psycopg.sql.Identifier( ProxyExtrasDBManager._prisma_schema_param(database_url) or "public" ), psycopg.sql.Identifier("_prisma_migrations"), ) try: with psycopg.connect( cleaned_url, connect_timeout=10, autocommit=True ) as conn: row = conn.execute( psycopg.sql.SQL( "SELECT logs FROM {} " "WHERE migration_name = %s AND finished_at IS NULL " "AND rolled_back_at IS NULL" ).format(ledger_table), (migration_name,), ).fetchone() except (psycopg.OperationalError, psycopg.DatabaseError): return None return (row[0] or "") if row else "" @staticmethod def _resolve_specific_migration(migration_name: str): """Mark a specific migration as applied""" prisma_env = _get_prisma_env() subprocess.run( [_get_prisma_command(), "migrate", "resolve", "--applied", migration_name], timeout=prisma_command_timeout(), check=True, capture_output=True, env=prisma_env, ) @staticmethod def _is_permission_error(error_message: str) -> bool: """ Check if the error message indicates a database permission error. Permission errors should NOT be marked as applied, as the migration did not actually execute successfully. Args: error_message: The error message from Prisma migrate Returns: bool: True if this is a permission error, False otherwise """ permission_patterns = [ r"Database error code: 42501", # PostgreSQL insufficient privilege r"must be owner of table", r"permission denied for schema", r"permission denied for table", r"must be owner of schema", ] for pattern in permission_patterns: if re.search(pattern, error_message, re.IGNORECASE): return True return False @staticmethod def _is_idempotent_error(error_message: str) -> bool: """ Check if the error message indicates an idempotent operation error. Idempotent errors (like "column already exists") mean the migration has effectively already been applied, so it's safe to mark as applied. Args: error_message: The error message from Prisma migrate Returns: bool: True if this is an idempotent error, False otherwise """ idempotent_patterns = [ r"already exists", r"column .* already exists", r"duplicate key value violates", r"relation .* already exists", r"constraint .* already exists", r"does not exist", r"Can't drop database.* because it doesn't exist", ] for pattern in idempotent_patterns: if re.search(pattern, error_message, re.IGNORECASE): return True return False @staticmethod def _resolve_all_migrations( migrations_dir: str, schema_path: str, mark_all_applied: bool = True ): """ 1. Compare the current database state to schema.prisma and generate a migration for the diff. 2. Run prisma migrate deploy to apply any pending migrations. 3. Mark all existing migrations as applied. """ database_url = os.getenv("DATABASE_URL") if not database_url: logger.error("DATABASE_URL not set") return # Prefer DIRECT_URL for schema introspection — pooler URLs (e.g. neon -pooler) # do not support the extended query protocol required by prisma migrate diff. diff_url = os.getenv("DIRECT_URL") or database_url diff_dir = Path(tempfile.mkdtemp(prefix="litellm_migration_diff_")) diff_sql_path = diff_dir / "migration.sql" # 1. Generate migration SQL for the diff between DB and schema try: logger.info("Generating migration diff between DB and schema.prisma...") with open(diff_sql_path, "w") as f: subprocess.run( [ _get_prisma_command(), "migrate", "diff", "--from-url", diff_url, "--to-schema-datamodel", schema_path, "--script", ], check=True, timeout=prisma_command_timeout(), stdout=f, env=_get_prisma_env(), ) except subprocess.CalledProcessError as e: logger.warning(f"Failed to generate migration diff: {e.stderr}") except subprocess.TimeoutExpired: logger.warning("Migration diff generation timed out.") # check if the migration was created if not diff_sql_path.exists(): logger.warning( "Migration diff was not created (prisma migrate diff failed — " "likely a pooler URL). Falling back to direct SQL execution of " "each migration file." ) # Fall back: run each migration SQL file directly via prisma db execute. # This works with pooler URLs (no schema introspection needed) and is # safe to re-run because migrations use IF NOT EXISTS / IF EXISTS guards. migration_files = sorted(Path(migrations_dir).glob("*/migration.sql")) for mig_file in migration_files: try: subprocess.run( [ _get_prisma_command(), "db", "execute", "--file", str(mig_file), "--schema", schema_path, ], timeout=prisma_command_timeout(), check=True, capture_output=True, text=True, env=_get_prisma_env(), ) logger.info(f"Applied migration: {mig_file.parent.name}") except subprocess.CalledProcessError as e: logger.warning( f"Failed to apply migration {mig_file.parent.name}: {e.stderr}" ) except subprocess.TimeoutExpired: logger.warning(f"Migration {mig_file.parent.name} timed out.") 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) return # 2. Run prisma db execute to apply the migration applied_ok = False try: logger.info("Running prisma db execute to apply the migration diff...") result = subprocess.run( [ _get_prisma_command(), "db", "execute", "--file", str(diff_sql_path), "--schema", schema_path, ], timeout=prisma_command_timeout(), check=True, capture_output=True, text=True, env=_get_prisma_env(), ) logger.info(f"prisma db execute stdout: {result.stdout}") logger.info("✅ Migration diff applied successfully") applied_ok = True except subprocess.CalledProcessError as e: logger.warning(f"Failed to apply migration diff: {e.stderr}") except subprocess.TimeoutExpired: logger.warning("Migration diff application timed out.") # 3. Mark all migrations as applied if not mark_all_applied: return if not applied_ok: logger.warning( "Drift script failed to apply; NOT marking migrations as " "applied so a later migration run can retry them" ) return ProxyExtrasDBManager._mark_migrations_applied(migrations_dir) @staticmethod def _mark_migrations_applied(migrations_dir: str) -> None: migration_names = ProxyExtrasDBManager._get_migration_names(migrations_dir) logger.info(f"Resolving {len(migration_names)} migrations") for migration_name in migration_names: try: logger.info(f"Resolving migration: {migration_name}") subprocess.run( [ _get_prisma_command(), "migrate", "resolve", "--applied", migration_name, ], timeout=prisma_command_timeout(), check=True, capture_output=True, text=True, env=_get_prisma_env(), ) logger.debug(f"Resolved migration: {migration_name}") except subprocess.CalledProcessError as e: if "is already recorded as applied in the database." not in e.stderr: logger.warning( f"Failed to resolve migration {migration_name}: {e.stderr}" ) @staticmethod def spend_logs_is_partitioned() -> bool: """True when the connected database's LiteLLM_SpendLogs is a partitioned table in Prisma's target schema (the `schema` URL param, falling back to Prisma's default target, public), i.e. the operator ran db_scripts/partition_spend_logs.sql. Returns False when psycopg is unavailable or the database cannot be reached, preserving the pre-existing behavior in those cases.""" database_url = os.getenv("DATABASE_URL") if not database_url: return False try: import psycopg except ImportError: logger.warning( "psycopg is not installed; skipping the LiteLLM_SpendLogs " "partition check. If this table is partitioned (see " "db_scripts/partition_spend_logs.sql), schema reconciliation " "will try to rewrite its primary key and fail. Install the " "litellm[extra_proxy] extra, which now includes psycopg." ) return False cleaned_url = ProxyExtrasDBManager._strip_prisma_query_params(database_url) try: with psycopg.connect( cleaned_url, connect_timeout=10, autocommit=True ) as conn: row = conn.execute( "SELECT 1 " "FROM pg_partitioned_table pt " "JOIN pg_class c ON c.oid = pt.partrelid " "JOIN pg_namespace n ON n.oid = c.relnamespace " "WHERE c.relname = 'LiteLLM_SpendLogs' " " AND n.nspname = %s", ( ProxyExtrasDBManager._prisma_schema_param(database_url) or "public", ), ).fetchone() except (psycopg.OperationalError, psycopg.DatabaseError): return False return row is not None @staticmethod def _prisma_schema_param(url: str) -> Optional[str]: """The `schema` query param Prisma uses to pick its target schema, or None when the URL does not set one.""" from urllib.parse import urlparse, parse_qsl return next( (v for k, v in parse_qsl(urlparse(url).query) if k == "schema"), None, ) @staticmethod def _strip_prisma_query_params(url: str) -> str: """Remove Prisma-specific query params (connection_limit, pool_timeout, schema, etc.) from DATABASE_URL so psycopg can parse it.""" from urllib.parse import urlparse, urlunparse, parse_qsl, urlencode parsed = urlparse(url) if not parsed.query: return url libpq_params = { "sslmode", "sslcert", "sslkey", "sslrootcert", "sslpassword", "application_name", "connect_timeout", "client_encoding", "options", "service", "gssencmode", "krbsrvname", "target_session_attrs", } kept = [(k, v) for k, v in parse_qsl(parsed.query) if k in libpq_params] return urlunparse(parsed._replace(query=urlencode(kept))) @staticmethod def _warn_if_db_ahead_of_head(migrations_dir: str) -> None: """ Log a warning if _prisma_migrations contains applied migrations with timestamps newer than every migration this build ships. This is informational only for the v2 resolver — it tells the operator the DB was likely migrated by a newer deployment, which is usually a signal that this (older) version shouldn't run against it. We do NOT block startup: many users have weird _prisma_migrations state from prior thrashing bugs, and blocking them would be a breaking change. Safe no-op if psycopg isn't installed or DB isn't reachable. """ database_url = os.getenv("DATABASE_URL") if not database_url: return try: import psycopg except ImportError: return cleaned_url = ProxyExtrasDBManager._strip_prisma_query_params(database_url) known = set(ProxyExtrasDBManager._get_migration_names(migrations_dir)) try: # autocommit=True keeps the SELECT outside a transaction. Without # it, psycopg3's `with conn` calls COMMIT on clean exit — which # fails after `UndefinedTable` (fresh DB) leaves the transaction # in an aborted state. with psycopg.connect( cleaned_url, connect_timeout=10, autocommit=True ) as conn: try: rows = conn.execute( "SELECT migration_name FROM _prisma_migrations " "WHERE finished_at IS NOT NULL AND rolled_back_at IS NULL" ).fetchall() except psycopg.errors.UndefinedTable: return except (psycopg.OperationalError, psycopg.DatabaseError): # Swallow connection failures AND any other DB-layer error # (e.g. InsufficientPrivilege if the runtime user lacks SELECT # on _prisma_migrations). This is an informational check — # never block startup on it. return applied = {r[0] for r in rows} unknown = applied - known if not unknown: return head_newest_ts = _max_migration_timestamp(known) hostile = { name for name in unknown if _migration_timestamp(name) > head_newest_ts } if not hostile: return sorted_hostile = sorted(hostile) logger.warning( "Database has %d migration(s) applied that are NEWER than any " "migration this LiteLLM version ships. This usually means the " "database was migrated by a newer LiteLLM deployment. Some API " "endpoints may fail because this proxy's Prisma client does not " "know about those schema changes. Consider upgrading this " "deployment. Unknown: %s", len(hostile), ", ".join(sorted_hostile[:5]) + (" ..." if len(sorted_hostile) > 5 else ""), ) @staticmethod def _setup_database_v2(use_migrate: bool) -> bool: """ v2 migration resolver (opt-in via --use_v2_migration_resolver). Runs `prisma migrate deploy` and handles standard recovery paths (P3005 baseline, P3009/P3018 idempotent errors, deadlocks against a concurrent migrate deploy). Critically, it does NOT call `_resolve_all_migrations` — the diff-and-force recovery that caused schema thrashing when two LiteLLM versions contended for the same DB during rolling deploys. Ahead-of-HEAD state (DB has migrations newer than this build ships) is logged as a warning, not a fatal error — users whose DBs got into weird shapes from the old thrashing should still be able to start. The retry budget only counts attempts that made no progress: see _MigrateAttemptBudget. """ schema_path = ProxyExtrasDBManager._get_prisma_dir() + "/schema.prisma" migrations_dir = ProxyExtrasDBManager._get_prisma_dir() if not use_migrate: if ProxyExtrasDBManager.spend_logs_is_partitioned(): raise RuntimeError(PARTITIONED_SPEND_LOGS_PUSH_ERROR) original_dir = os.getcwd() os.chdir(migrations_dir) try: subprocess.run( [_get_prisma_command(), "db", "push", "--accept-data-loss"], timeout=prisma_command_timeout(), check=True, env=_get_prisma_env(), ) return True except ( subprocess.CalledProcessError, subprocess.TimeoutExpired, ) as e: # Re-raise as RuntimeError so proxy_cli.py's # `except RuntimeError` catches it and exits cleanly. raise RuntimeError(f"prisma db push failed.\n\nDetail: {e}") from e finally: os.chdir(original_dir) # Informational — never blocks. ProxyExtrasDBManager._warn_if_db_ahead_of_head(migrations_dir) original_dir = os.getcwd() os.chdir(migrations_dir) deploy_timeout = prisma_migrate_deploy_timeout() budget = _MigrateAttemptBudget(attempts_left=MAX_MIGRATE_DEPLOY_ATTEMPTS) try: while not budget.exhausted: try: result = subprocess.run( [_get_prisma_command(), "migrate", "deploy"], timeout=deploy_timeout, check=True, capture_output=True, text=True, env=_get_prisma_env(), ) logger.info(f"prisma migrate deploy stdout: {result.stdout}") return True except subprocess.TimeoutExpired: logger.warning( "prisma migrate deploy attempt %s timed out after %ss, retrying. " "Raise %s if this database needs longer to apply its pending migrations.", budget.attempt_number, deploy_timeout, PRISMA_MIGRATE_DEPLOY_TIMEOUT_ENV_VAR, ) next_budget = budget.spend() except subprocess.CalledProcessError as e: next_budget = ProxyExtrasDBManager._budget_after_deploy_failure( e, budget, schema_path ) if next_budget.attempts_left < budget.attempts_left: time.sleep(random.randrange(5, 15)) budget = next_budget # rebind-ok: the loop carries the budget from one migrate deploy pass to the next raise RuntimeError( f"Database migration failed after {MAX_MIGRATE_DEPLOY_ATTEMPTS} " "attempts that made no progress (timeouts, deadlock retries, or a " "recovery that had already run once). Check database connectivity, " "load, and _prisma_migrations ledger state, and raise " f"{PRISMA_MIGRATE_DEPLOY_TIMEOUT_ENV_VAR} if the attempts timed out." ) finally: os.chdir(original_dir) @staticmethod def _budget_after_deploy_failure( error: subprocess.CalledProcessError, budget: "_MigrateAttemptBudget", schema_path: str, ) -> "_MigrateAttemptBudget": """Recover from one failed `prisma migrate deploy`, and price the pass. Returns the budget the next pass runs under, or raises when the failure is not one this resolver knows how to recover from. """ stderr = error.stderr or "" if "P3005" in stderr and "database schema is not empty" in stderr: logger.info("Schema exists but no migrations ledger — creating baseline") if ProxyExtrasDBManager._create_baseline_migration(schema_path): return budget.after_recovery("baseline") return budget.spend() if "P3009" in stderr: migration_match = re.search(r"`(\d+_\S+?)`", stderr) if migration_match and ProxyExtrasDBManager._is_idempotent_error(stderr): name = migration_match.group(1) logger.info( f"Migration {name} failed idempotently — marking applied and retrying" ) ProxyExtrasDBManager._mark_migration_applied(name) return budget.after_recovery(f"resolved:{name}") if migration_match: migration_name = migration_match.group(1) ledger_logs = ProxyExtrasDBManager._failed_migration_logs(migration_name) if ledger_logs is not None and ( ledger_logs == "" or _MIGRATION_DEADLOCK_MARKER in ledger_logs ): logger.info( "Migration %s failed in a concurrent migrate deploy " "deadlock race, rolling its ledger row back and retrying", migration_name, ) ProxyExtrasDBManager._roll_back_migration_best_effort(migration_name) return budget.spend() raise RuntimeError( "Database migration failed and cannot be auto-recovered. " f"Manual intervention required.\n\nPrisma error:\n{stderr}" ) from error if "P3018" in stderr: if ProxyExtrasDBManager._is_permission_error(stderr): raise RuntimeError( "Database migration failed due to insufficient " "permissions. Please grant the required privileges " f"and retry.\n\nPrisma error:\n{stderr}" ) from error migration_match = re.search(r"Migration name: (\d+_\S+)", stderr) if migration_match and ProxyExtrasDBManager._is_idempotent_error(stderr): name = migration_match.group(1) logger.info( f"Migration {name} SQL hit idempotent error — marking applied and retrying" ) ProxyExtrasDBManager._mark_migration_applied(name) return budget.after_recovery(f"resolved:{name}") if migration_match and _MIGRATION_DEADLOCK_MARKER in stderr: logger.info( "Migration %s deadlocked against a concurrent " "migrate deploy, rolling its ledger row back " "and retrying", migration_match.group(1), ) ProxyExtrasDBManager._roll_back_migration_best_effort( migration_match.group(1) ) return budget.spend() raise RuntimeError( "Database migration failed and cannot be auto-recovered. " f"Manual intervention required.\n\nPrisma error:\n{stderr}" ) from error if _MIGRATION_DEADLOCK_MARKER in stderr: logger.info( "prisma migrate deploy attempt %s deadlocked against " "a concurrent migrate deploy, retrying", budget.attempt_number, ) return budget.spend() if "P1002" in stderr and "advisory lock" in stderr: logger.info( "prisma migrate deploy attempt %s timed out waiting for " "the advisory lock a concurrent migrate deploy holds, retrying", budget.attempt_number, ) return budget.spend() raise RuntimeError( "Database migration failed and cannot be auto-recovered. " f"Manual intervention required.\n\nPrisma error:\n{stderr}" ) from error @staticmethod def _mark_migration_applied(name: str) -> None: """Roll a failed ledger row back if it is still there, then mark it applied.""" try: ProxyExtrasDBManager._roll_back_migration(name) except (subprocess.CalledProcessError, subprocess.TimeoutExpired): pass # may already be rolled-back try: ProxyExtrasDBManager._resolve_specific_migration(name) except ( subprocess.CalledProcessError, subprocess.TimeoutExpired, ) as resolve_err: # We're called from inside an `except CalledProcessError` handler — # re-raising CalledProcessError from here would escape as itself, # bypassing proxy_cli.py's `except RuntimeError`. raise RuntimeError( f"Failed to mark migration {name} as applied " f"after idempotent recovery. Manual " f"intervention may be required.\n\n" f"Detail: {resolve_err}" ) from resolve_err @staticmethod def apply_replica_identity_full_if_requested() -> bool: """ Re-assert REPLICA IDENTITY FULL on LiteLLM's tables when the operator opted in via LITELLM_SET_REPLICA_IDENTITY_FULL. Prisma leaves new tables at the Postgres default, which logical replication consumers reject, so the setting has to be re-applied after every migration run rather than once by hand. Returns: bool: True if the setting was applied, False if it was not requested or could not be applied. """ if not str_to_bool(os.getenv(REPLICA_IDENTITY_FULL_ENV_VAR)): return False try: schema_path = ProxyExtrasDBManager._get_prisma_dir() + "/schema.prisma" prisma_command = _get_prisma_command() prisma_env = _get_prisma_env() except OSError as e: logger.error( "Could not resolve the migrations directory for the REPLICA " "IDENTITY FULL step, skipping it. Error: %s", e, ) return False return apply_replica_identity_full( schema_path=schema_path, prisma_command=prisma_command, prisma_env=prisma_env, ) @staticmethod 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 Args: use_migrate: Whether to use prisma migrate instead of db push use_v2_resolver: Opt into the v2 migration resolver (safer during rolling deploys; does not run the diff-and-force recovery that causes schema thrashing). Defaults to False for backwards compatibility. Returns: bool: True if setup was successful, False otherwise """ ensure_prisma_toolchain( prisma_command=_get_prisma_command(), prisma_env=_get_prisma_env() ) migrated = ProxyExtrasDBManager._run_migrations( use_migrate=use_migrate, use_v2_resolver=use_v2_resolver ) if migrated: ProxyExtrasDBManager.apply_replica_identity_full_if_requested() return migrated @staticmethod def _run_migrations(use_migrate: bool, use_v2_resolver: bool) -> bool: if use_v2_resolver: logger.info("Using v2 migration resolver (--use_v2_migration_resolver)") return ProxyExtrasDBManager._setup_database_v2(use_migrate=use_migrate) schema_path = ProxyExtrasDBManager._get_prisma_dir() + "/schema.prisma" for attempt in range(4): original_dir = os.getcwd() migrations_dir = ProxyExtrasDBManager._get_prisma_dir() os.chdir(migrations_dir) try: if use_migrate: logger.info("Running prisma migrate deploy") try: # Set migrations directory for Prisma result = subprocess.run( [_get_prisma_command(), "migrate", "deploy"], timeout=prisma_migrate_deploy_timeout(), check=True, capture_output=True, text=True, env=_get_prisma_env(), ) logger.info(f"prisma migrate deploy stdout: {result.stdout}") logger.info("prisma migrate deploy completed") # Skip sanity check when deploy reports no pending migrations — # DB already matches schema, no drift to correct. if "No pending migrations to apply" in result.stdout: logger.info( "No pending migrations — skipping post-migration sanity check" ) return True # Run sanity check to ensure DB matches schema logger.info("Running post-migration sanity check...") ProxyExtrasDBManager._resolve_all_migrations( migrations_dir, schema_path, mark_all_applied=False ) 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: # Extract the failed migration name from the error message migration_match = re.search( r"`(\d+_.*)` migration", e.stderr ) if migration_match: failed_migration = migration_match.group(1) if ProxyExtrasDBManager._is_idempotent_error(e.stderr): logger.info( f"Migration {failed_migration} failed due to idempotent error (e.g., column already exists), resolving as applied" ) try: ProxyExtrasDBManager._roll_back_migration( failed_migration ) except ( subprocess.CalledProcessError, subprocess.TimeoutExpired, ) as rollback_err: logger.warning( f"Failed to roll back migration {failed_migration}: {rollback_err}. " f"It may already be in a rolled-back state." ) try: ProxyExtrasDBManager._resolve_specific_migration( failed_migration ) logger.info( f"✅ Migration {failed_migration} resolved, retrying to apply remaining migrations" ) except ( subprocess.CalledProcessError, subprocess.TimeoutExpired, ) as resolve_err: logger.warning( f"Failed to resolve migration {failed_migration}: {resolve_err}" ) # Apply any schema drift not covered by the marked-as-applied migration ProxyExtrasDBManager._resolve_all_migrations( migrations_dir, schema_path, mark_all_applied=False, ) else: logger.info( f"Found failed migration: {failed_migration}, marking as rolled back" ) # Mark the failed migration as rolled back subprocess.run( [ _get_prisma_command(), "migrate", "resolve", "--rolled-back", failed_migration, ], timeout=prisma_command_timeout(), check=True, capture_output=True, text=True, env=_get_prisma_env(), ) logger.info( f"✅ Migration {failed_migration} marked as rolled back... retrying" ) elif ( "P3005" in e.stderr and "database schema is not empty" in e.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" ) ProxyExtrasDBManager._create_baseline_migration(schema_path) logger.info( "Baseline migration created, resolving all migrations" ) ProxyExtrasDBManager._resolve_all_migrations( migrations_dir, schema_path ) logger.info("✅ All migrations resolved.") return True elif "P3018" in e.stderr: # Check if this is a permission error or idempotent error if ProxyExtrasDBManager._is_permission_error(e.stderr): # Permission errors should NOT be marked as applied # Extract migration name for logging migration_match = re.search( r"Migration name: (\d+_.*)", e.stderr ) migration_name = ( migration_match.group(1) if migration_match else "unknown" ) logger.error( f"❌ Migration {migration_name} failed due to insufficient permissions. " f"Please check database user privileges. Error: {e.stderr}" ) # Mark as rolled back and exit with error if migration_match: try: ProxyExtrasDBManager._roll_back_migration( migration_name ) logger.info( f"Migration {migration_name} marked as rolled back" ) except Exception as rollback_error: logger.warning( f"Failed to mark migration as rolled back: {rollback_error}" ) # Re-raise the error to prevent silent failures raise RuntimeError( f"Migration failed due to permission error. Migration {migration_name} " f"was NOT applied. Please grant necessary database permissions and retry." ) from e elif ProxyExtrasDBManager._is_idempotent_error(e.stderr): # Idempotent errors mean the migration has effectively been applied logger.info( "Migration failed due to idempotent error (e.g., column already exists), " "resolving as applied" ) # Extract the migration name from the error message migration_match = re.search( r"Migration name: (\d+_.*)", e.stderr ) if migration_match: migration_name = migration_match.group(1) try: logger.info( f"Rolling back migration {migration_name}" ) ProxyExtrasDBManager._roll_back_migration( migration_name ) except ( subprocess.CalledProcessError, subprocess.TimeoutExpired, ) as rollback_err: logger.warning( f"Failed to roll back migration {migration_name}: {rollback_err}. " f"It may already be in a rolled-back state." ) try: logger.info( f"Resolving migration {migration_name} that failed " f"due to existing schema objects" ) ProxyExtrasDBManager._resolve_specific_migration( migration_name ) logger.info( f"✅ Migration {migration_name} resolved, " f"retrying to apply remaining migrations" ) except ( subprocess.CalledProcessError, subprocess.TimeoutExpired, ) as resolve_err: logger.warning( f"Failed to resolve migration {migration_name}: {resolve_err}" ) # Apply any schema drift not covered by the marked-as-applied migration ProxyExtrasDBManager._resolve_all_migrations( migrations_dir, schema_path, mark_all_applied=False, ) else: # Unknown P3018 error - log and re-raise for safety logger.warning( f"P3018 error encountered but could not classify " f"as permission or idempotent error. " f"Error: {e.stderr}" ) raise else: if ProxyExtrasDBManager.spend_logs_is_partitioned(): raise RuntimeError(PARTITIONED_SPEND_LOGS_PUSH_ERROR) # Use prisma db push with increased timeout subprocess.run( [_get_prisma_command(), "db", "push", "--accept-data-loss"], timeout=prisma_command_timeout(), check=True, ) return True except subprocess.TimeoutExpired: logger.warning( "Attempt %s timed out. Raise %s if this database needs longer to apply its schema.", attempt + 1, PRISMA_MIGRATE_DEPLOY_TIMEOUT_ENV_VAR if use_migrate else PRISMA_COMMAND_TIMEOUT_ENV_VAR, ) time.sleep(random.randrange(5, 15)) except subprocess.CalledProcessError as e: attempts_left = 3 - attempt retry_msg = ( f" Retrying... ({attempts_left} attempts left)" if attempts_left > 0 else "" ) logger.info(f"The process failed to execute. Details: {e}.{retry_msg}") time.sleep(random.randrange(5, 15)) finally: os.chdir(original_dir) pass return False