diff --git a/litellm-proxy-extras/litellm_proxy_extras/migration_lock.py b/litellm-proxy-extras/litellm_proxy_extras/migration_lock.py index bea5e36fd18..3df82d8c81e 100644 --- a/litellm-proxy-extras/litellm_proxy_extras/migration_lock.py +++ b/litellm-proxy-extras/litellm_proxy_extras/migration_lock.py @@ -1,4 +1,5 @@ import random +import re import time from collections.abc import Generator, Mapping from contextlib import contextmanager @@ -7,25 +8,50 @@ from typing import TYPE_CHECKING, Final from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit from litellm_proxy_extras._logging import logger -from litellm_proxy_extras.prisma_toolchain import MIGRATION_LOCK_TIMEOUT_ENV_VAR, migration_lock_timeout +from litellm_proxy_extras.prisma_toolchain import ( + MIGRATION_LOCK_TIMEOUT_ENV_VAR, + migration_ddl_lock_timeout, + migration_lock_timeout, +) MIGRATION_LOCK_KEY: Final = int.from_bytes(b"llm_mig2", "big") +_LOCK_TIMEOUT_PINNED_RE: Final = re.compile(r"(?:-c\s*|--)lock_timeout=") if TYPE_CHECKING: import psycopg -def migration_environment(environment: Mapping[str, str]) -> Mapping[str, str]: - database_url: Final = environment.get("DATABASE_URL") - direct_url: Final = environment.get("DIRECT_URL") - if not database_url or not direct_url: - return environment +def _migration_url(database_url: str, direct_url: "str | None") -> str: + if not direct_url: + return database_url schema: Final = next((value for key, value in parse_qsl(urlsplit(database_url).query) if key == "schema"), "public") direct: Final = urlsplit(direct_url) parameters: Final = tuple((key, value) for key, value in parse_qsl(direct.query) if key != "schema") + return urlunsplit(direct._replace(query=urlencode((*parameters, ("schema", schema))))) + + +def _with_ddl_lock_timeout(url: str) -> str: + split: Final = urlsplit(url) + pairs: Final = tuple(parse_qsl(split.query)) + if any(key == "pgbouncer" and value == "true" for key, value in pairs): + return url + existing: Final = next((value for key, value in pairs if key == "options"), "") + if _LOCK_TIMEOUT_PINNED_RE.search(existing): + return url + timeout_ms: Final = max(1, int(migration_ddl_lock_timeout() * 1000)) + options: Final = f"{existing} -c lock_timeout={timeout_ms}".strip() + rewritten: Final = tuple((key, value) for key, value in pairs if key != "options") + (("options", options),) + return urlunsplit(split._replace(query=urlencode(rewritten))) + + +def migration_environment(environment: Mapping[str, str]) -> Mapping[str, str]: + database_url: Final = environment.get("DATABASE_URL") + if not database_url: + return environment + migration_url: Final = _migration_url(database_url, environment.get("DIRECT_URL")) return { **environment, - "DATABASE_URL": urlunsplit(direct._replace(query=urlencode((*parameters, ("schema", schema))))), + "DATABASE_URL": _with_ddl_lock_timeout(migration_url), } diff --git a/litellm-proxy-extras/litellm_proxy_extras/prisma_toolchain.py b/litellm-proxy-extras/litellm_proxy_extras/prisma_toolchain.py index 07f83f76d2b..11b58284f41 100644 --- a/litellm-proxy-extras/litellm_proxy_extras/prisma_toolchain.py +++ b/litellm-proxy-extras/litellm_proxy_extras/prisma_toolchain.py @@ -60,11 +60,13 @@ PRISMA_COMMAND_TIMEOUT_ENV_VAR = "LITELLM_PRISMA_COMMAND_TIMEOUT" PRISMA_BOOTSTRAP_TIMEOUT_ENV_VAR = "LITELLM_PRISMA_BOOTSTRAP_TIMEOUT" PRISMA_MIGRATE_DEPLOY_TIMEOUT_ENV_VAR = "LITELLM_PRISMA_MIGRATE_DEPLOY_TIMEOUT" MIGRATION_LOCK_TIMEOUT_ENV_VAR = "LITELLM_MIGRATION_LOCK_TIMEOUT" +MIGRATION_DDL_LOCK_TIMEOUT_ENV_VAR = "LITELLM_MIGRATION_DDL_LOCK_TIMEOUT" NODEENV_CACHE_DIR_ENV_VAR = "PRISMA_NODEENV_CACHE_DIR" DEFAULT_PRISMA_COMMAND_TIMEOUT = 60.0 DEFAULT_PRISMA_BOOTSTRAP_TIMEOUT = 600.0 DEFAULT_PRISMA_MIGRATE_DEPLOY_TIMEOUT = 600.0 +DEFAULT_MIGRATION_DDL_LOCK_TIMEOUT = 10.0 BOOTSTRAP_ARG = "--version" PRISMA_CONSOLE_SCRIPT = "prisma" @@ -111,6 +113,11 @@ def migration_lock_timeout() -> float: return _timeout_from_env(MIGRATION_LOCK_TIMEOUT_ENV_VAR, 600.0) +def migration_ddl_lock_timeout() -> float: + """Seconds one migration statement may wait on a table lock before Postgres cancels it.""" + return _timeout_from_env(MIGRATION_DDL_LOCK_TIMEOUT_ENV_VAR, DEFAULT_MIGRATION_DDL_LOCK_TIMEOUT) + + def prisma_bootstrap_timeout() -> float: """Seconds the one-time Node toolchain install may run for.""" return _timeout_from_env( diff --git a/litellm-proxy-extras/litellm_proxy_extras/utils.py b/litellm-proxy-extras/litellm_proxy_extras/utils.py index 72462fa9ae4..061fb627c65 100644 --- a/litellm-proxy-extras/litellm_proxy_extras/utils.py +++ b/litellm-proxy-extras/litellm_proxy_extras/utils.py @@ -18,6 +18,7 @@ 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 ( + MIGRATION_DDL_LOCK_TIMEOUT_ENV_VAR, PRISMA_COMMAND_TIMEOUT_ENV_VAR, PRISMA_MIGRATE_DEPLOY_TIMEOUT_ENV_VAR, ensure_prisma_toolchain, @@ -62,6 +63,7 @@ def _get_prisma_env() -> dict: _MIGRATION_TS_RE = re.compile(r"^(\d{14})_") _MIGRATION_DEADLOCK_MARKER = "deadlock detected" +_MIGRATION_LOCK_TIMEOUT_MARKER = "canceling statement due to lock timeout" INDEX_REPAIR_ADVISORY_LOCK_KEY: Final = int.from_bytes(b"litellm", "big") _TRANSIENT_INDEX_SUFFIX_RE: Final = re.compile(r"_cc(?:new|old)\d*$") _INVALID_LITELLM_INDEXES_SQL: Final = ( @@ -1162,9 +1164,11 @@ class ProxyExtrasDBManager: raise RuntimeError( f"Database migration failed after {MAX_MIGRATE_DEPLOY_ATTEMPTS} " - "attempts that made no progress (timeouts or deadlock retries). Check database connectivity, " - "load, and _prisma_migrations ledger state, and raise " - f"{PRISMA_MIGRATE_DEPLOY_TIMEOUT_ENV_VAR} if the attempts timed out." + "attempts that made no progress (timeouts, lock timeouts or deadlock retries). " + "Check database connectivity, load, long-running transactions on the migrated tables, " + "and _prisma_migrations ledger state. Raise " + f"{PRISMA_MIGRATE_DEPLOY_TIMEOUT_ENV_VAR} if the attempts timed out, or " + f"{MIGRATION_DDL_LOCK_TIMEOUT_ENV_VAR} if they timed out waiting for a table lock." ) finally: os.chdir(original_dir) @@ -1225,6 +1229,13 @@ class ProxyExtrasDBManager: ) ProxyExtrasDBManager._v2_roll_back_migration_best_effort(migration_name) return budget.spend() + if ledger_logs and _MIGRATION_LOCK_TIMEOUT_MARKER in ledger_logs: + logger.info( + "Migration %s timed out waiting for a table lock, rolling its ledger row back and retrying", + migration_name, + ) + ProxyExtrasDBManager._v2_roll_back_migration_best_effort(migration_name) + return budget.spend() if ProxyExtrasDBManager._failed_migration_recovered(migration_name, started_at): logger.info( "Migration %s started at %s was already rolled back or completed by a concurrent " @@ -1270,6 +1281,16 @@ class ProxyExtrasDBManager: ProxyExtrasDBManager._v2_roll_back_migration_best_effort(migration_name) return budget.spend() + if migration_name and _MIGRATION_LOCK_TIMEOUT_MARKER in stderr: + ProxyExtrasDBManager._log_migration_lock_holders(migration_name) + logger.warning( + "Migration %s timed out waiting for a table lock, rolling its ledger row back and retrying. " + f"Raise {MIGRATION_DDL_LOCK_TIMEOUT_ENV_VAR} if the database needs longer.", + migration_name, + ) + ProxyExtrasDBManager._v2_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}" @@ -1282,6 +1303,13 @@ class ProxyExtrasDBManager: ) return budget.spend() + if _MIGRATION_LOCK_TIMEOUT_MARKER in stderr: + logger.info( + "Waiting for the advisory lock held by another Prisma migration; " + "contention does not spend a migration failure attempt" + ) + return budget.after_contention(attempt_seconds) + if "P1002" in stderr and "advisory lock" in stderr: logger.info( "Waiting for the advisory lock held by another Prisma migration; " @@ -1294,6 +1322,68 @@ class ProxyExtrasDBManager: f"Manual intervention required.\n\nPrisma error:\n{stderr}" ) from error + @staticmethod + def _log_migration_lock_holders(migration_name: str) -> None: + """Best-effort log of the sessions holding locks on the tables a timed-out migration touches.""" + try: + migration_sql: Final = ( + Path(ProxyExtrasDBManager._get_prisma_dir()) / "migrations" / migration_name / "migration.sql" + ).read_text() + except OSError: + return + relations: Final = sorted(set(re.findall(r'"(LiteLLM_\w+)"', migration_sql))) + database_url: Final = os.getenv("DATABASE_URL") + if not relations or not database_url: + return + try: + import psycopg + except ImportError: + return + try: + with psycopg.connect( + ProxyExtrasDBManager._strip_prisma_query_params(database_url), + connect_timeout=10, + autocommit=True, + options="-c statement_timeout=5000", + ) as conn: + rows: Final = conn.execute( + "SELECT l.pid, c.relname, l.mode, a.state, a.application_name, " + "date_trunc('second', now() - a.xact_start) " + "FROM pg_locks l " + "JOIN pg_class c ON c.oid = l.relation " + "JOIN pg_namespace n ON n.oid = c.relnamespace " + "JOIN pg_stat_activity a ON a.pid = l.pid " + "WHERE l.granted AND n.nspname = %s AND c.relname = ANY(%s) AND l.pid <> pg_backend_pid() " + "ORDER BY a.xact_start NULLS LAST " + "LIMIT 10", + ( + ProxyExtrasDBManager._prisma_schema_param(database_url) or "public", + relations, + ), + ).fetchall() + except psycopg.Error: + return + if not rows: + logger.warning( + "Migration %s timed out waiting for a lock, but no current lock holder " + "was found on %s; the holder has likely since committed or rolled back", + migration_name, + ", ".join(relations), + ) + return + for pid, relname, mode, state, application_name, age in rows: + logger.warning( + "Migration %s timed out waiting for a lock on %s held by pid %s " + "(mode %s, state %s, application %s, transaction open for %s)", + migration_name, + relname, + pid, + mode, + state, + application_name, + age, + ) + @staticmethod def _mark_migration_applied(name: str) -> None: """Roll a failed ledger row back if it is still there, then mark it applied.""" diff --git a/tests/unit/litellm_proxy_extras/test_litellm_proxy_extras_utils.py b/tests/unit/litellm_proxy_extras/test_litellm_proxy_extras_utils.py index 7edb531a63b..a4ed691f5ee 100644 --- a/tests/unit/litellm_proxy_extras/test_litellm_proxy_extras_utils.py +++ b/tests/unit/litellm_proxy_extras/test_litellm_proxy_extras_utils.py @@ -766,6 +766,28 @@ _P3005_STDERR = """Error: P3005 The database schema is not empty. Read more about how to baseline an existing production database: https://pris.ly/d/migrate-baseline """ +_LOCK_TIMEOUT_MIGRATION = "20260818000000_add_spend_log_timestamps" + +_P3018_LOCK_TIMEOUT_STDERR = ( + "Error: P3018\n\n" + "A migration failed to apply. New migrations cannot be applied before the error is recovered from. " + "Read more about how to resolve migration issues in a production database: " + "https://pris.ly/d/migrate-resolve\n\n" + "Migration name: 20260818000000_add_spend_log_timestamps\n\n" + "Database error code: 55P03\n\n" + "Database error:\n" + "ERROR: canceling statement due to lock timeout\n\n" + 'DbError { severity: "ERROR", parsed_severity: Some(Error), code: SqlState(E55P03), ' + 'message: "canceling statement due to lock timeout", detail: None, hint: None, position: None, ' + "where_: None, schema: None, table: None, column: None, datatype: None, constraint: None, " + 'file: Some("postgres.c"), line: Some(3298), routine: Some("ProcessInterrupts") }\n' +) + +_LOCK_TIMEOUT_CONTENTION_STDERR = """Error: ERROR: canceling statement due to lock timeout + 0: schema_core::state::ApplyMigrations + at schema-engine/core/src/state.rs:201 +""" + def _p3018_stderr(migration_name): return f"""Error: P3018 @@ -859,17 +881,22 @@ class _FakeLedger: "pooled,direct,expected", ( ("postgresql://pool/db?pgbouncer=true", None, "postgresql://pool/db?pgbouncer=true"), - ("postgresql://pool/db?pgbouncer=true", "postgresql://writer/db", "postgresql://writer/db?schema=public"), + ( + "postgresql://pool/db?pgbouncer=true", + "postgresql://writer/db", + "postgresql://writer/db?schema=public&options=-c+lock_timeout%3D10000", + ), ( "postgresql://pool/db?schema=tenant%20one&pgbouncer=true", "postgresql://writer/db?sslmode=require&schema=wrong", - "postgresql://writer/db?sslmode=require&schema=tenant+one", + "postgresql://writer/db?sslmode=require&schema=tenant+one&options=-c+lock_timeout%3D10000", ), ), ) -def test_v2_migrations_use_the_direct_connection_with_the_runtime_schema(pooled, direct, expected): +def test_v2_migrations_use_the_direct_connection_with_the_runtime_schema(monkeypatch, pooled, direct, expected): from litellm_proxy_extras.migration_lock import migration_environment + monkeypatch.delenv("LITELLM_MIGRATION_DDL_LOCK_TIMEOUT", raising=False) environment = {"DATABASE_URL": pooled, "PRISMA_OFFLINE_MODE": "true"} configured = {**environment, **({"DIRECT_URL": direct} if direct else {})} migrated = migration_environment(configured) @@ -879,6 +906,106 @@ def test_v2_migrations_use_the_direct_connection_with_the_runtime_schema(pooled, assert configured["DATABASE_URL"] == pooled +@pytest.mark.parametrize( + "database_url,direct_url,ddl_timeout,expected_pairs", + ( + ( + "postgresql://writer/db", + None, + None, + [("options", "-c lock_timeout=10000")], + ), + ( + "postgresql://pool/db?schema=tenant&pgbouncer=true", + "postgresql://writer/db?options=-c%20statement_timeout%3D0&schema=wrong", + None, + [("schema", "tenant"), ("options", "-c statement_timeout=0 -c lock_timeout=10000")], + ), + ( + "postgresql://pool/db", + "postgresql://writer/db?options=-c%20lock_timeout%3D2500", + None, + [("options", "-c lock_timeout=2500"), ("schema", "public")], + ), + ( + "postgresql://pool/db", + "postgresql://writer/db?options=-clock_timeout%3D2500", + None, + [("options", "-clock_timeout=2500"), ("schema", "public")], + ), + ( + "postgresql://pool/db", + "postgresql://writer/db?options=--lock_timeout%3D0", + None, + [("options", "--lock_timeout=0"), ("schema", "public")], + ), + ( + "postgresql://pool/db", + "postgresql://writer/db?options=-c%20deadlock_timeout%3D1000", + None, + [("schema", "public"), ("options", "-c deadlock_timeout=1000 -c lock_timeout=10000")], + ), + ( + "postgresql://pool/db", + "postgresql://writer/db?pgbouncer=true", + None, + [("pgbouncer", "true"), ("schema", "public")], + ), + ( + "postgresql://writer/db", + None, + "2.5", + [("options", "-c lock_timeout=2500")], + ), + ( + "postgresql://writer/db", + None, + "abc", + [("options", "-c lock_timeout=10000")], + ), + ), + ids=( + "plain-database-url", + "direct-url-with-an-existing-options-value", + "pinned-dash-c-lock-timeout", + "pinned-dash-clock-timeout", + "pinned-double-dash-lock-timeout-zero-disables", + "unrelated-deadlock-timeout-is-not-a-pin", + "pooled-direct-url", + "timeout-override", + "invalid-override-falls-back", + ), +) +def test_migration_environment_pins_a_ddl_lock_timeout( + monkeypatch, database_url, direct_url, ddl_timeout, expected_pairs +): + from urllib.parse import parse_qsl, urlsplit + + from litellm_proxy_extras.migration_lock import migration_environment + + if ddl_timeout is None: + monkeypatch.delenv("LITELLM_MIGRATION_DDL_LOCK_TIMEOUT", raising=False) + else: + monkeypatch.setenv("LITELLM_MIGRATION_DDL_LOCK_TIMEOUT", ddl_timeout) + environment = {"DATABASE_URL": database_url, **({"DIRECT_URL": direct_url} if direct_url else {})} + snapshot = dict(environment) + + migrated = migration_environment(environment) + + assert parse_qsl(urlsplit(migrated["DATABASE_URL"]).query) == expected_pairs + assert environment == snapshot + + +def test_migration_environment_lock_timeout_encodes_the_options_value(monkeypatch): + from litellm_proxy_extras.migration_lock import migration_environment + + monkeypatch.delenv("LITELLM_MIGRATION_DDL_LOCK_TIMEOUT", raising=False) + + migrated = migration_environment({"DATABASE_URL": "postgresql://writer/db"}) + + assert migrated["DATABASE_URL"] == "postgresql://writer/db?options=-c+lock_timeout%3D10000" + + class _MigrateDeployHarness: """Drives _setup_database_v2 with a scripted sequence of `prisma migrate deploy` outcomes, with every recovery command faked out so @@ -899,6 +1026,8 @@ class _MigrateDeployHarness: self.deploy_calls = [] self.resolved = [] + self.rolled_back = [] + self.sleeps = [] self.baselines = 0 self._outcomes = list(outcomes) self._repeat_last = repeat_last @@ -913,7 +1042,7 @@ class _MigrateDeployHarness: monkeypatch.setenv("LITELLM_MIGRATION_DIR", str(tmp_path)) monkeypatch.setattr(utils_module.prisma_toolchain, "run_prisma", self._fake_run) monkeypatch.setattr(utils_module, "_get_prisma_env", lambda: {}) - monkeypatch.setattr(utils_module.time, "sleep", lambda seconds: None) + monkeypatch.setattr(utils_module.time, "sleep", self.sleeps.append) self.baseline_succeeds = True @@ -930,14 +1059,18 @@ class _MigrateDeployHarness: raise AssertionError("prisma migrate deploy called more times than scripted") def _fake_run(self, cmd, **kwargs): - assert cmd[1:] == ["migrate", "deploy"], f"unexpected prisma command: {cmd}" - self.deploy_calls.append(cmd) - outcome = self._next_outcome() - if outcome == "ok": + if cmd[1:] == ["migrate", "deploy"]: + self.deploy_calls.append(cmd) + outcome = self._next_outcome() + if outcome == "ok": + return _FakeCompleted() + if outcome == "timeout": + raise self._subprocess_module.TimeoutExpired(cmd, 1) + raise self._subprocess_module.CalledProcessError(1, cmd, stderr=outcome) + if cmd[1:4] == ["migrate", "resolve", "--rolled-back"] and len(cmd) == 5: + self.rolled_back.append(cmd[4]) return _FakeCompleted() - if outcome == "timeout": - raise self._subprocess_module.TimeoutExpired(cmd, 1) - raise self._subprocess_module.CalledProcessError(1, cmd, stderr=outcome) + raise AssertionError(f"unexpected prisma command: {cmd}") def run(self): while not ProxyExtrasDBManager._run_database_v2( @@ -1147,6 +1280,156 @@ class TestMigrateDeployAttemptAccounting: assert harness.resolved == [] +class TestMigrationLockTimeoutRecovery: + """A migration whose DDL wait Postgres cancelled (55P03, from the lock_timeout the + migration URL now pins) is retryable: the half-written ledger row is rolled back and + the next attempt reruns the migration, priced like a deadlock retry.""" + + def test_a_lock_timed_out_migration_rolls_back_and_retries(self, monkeypatch, tmp_path): + harness = _MigrateDeployHarness( + monkeypatch, + tmp_path, + [_P3018_LOCK_TIMEOUT_STDERR, "ok"], + ) + + assert harness.run() is True + assert len(harness.deploy_calls) == 2 + assert harness.rolled_back == [_LOCK_TIMEOUT_MIGRATION] + assert len(harness.sleeps) == 1 + + def test_repeated_lock_timeouts_spend_the_failure_budget(self, monkeypatch, tmp_path): + harness = _MigrateDeployHarness( + monkeypatch, + tmp_path, + [_P3018_LOCK_TIMEOUT_STDERR], + repeat_last=True, + ) + + with pytest.raises(RuntimeError, match="LITELLM_MIGRATION_DDL_LOCK_TIMEOUT"): + harness.run() + assert len(harness.deploy_calls) == _ATTEMPT_BUDGET + assert harness.rolled_back == [_LOCK_TIMEOUT_MIGRATION] * _ATTEMPT_BUDGET + + def test_lock_timeout_on_prismas_own_lock_wait_is_contention(self, monkeypatch, tmp_path): + harness = _MigrateDeployHarness( + monkeypatch, + tmp_path, + [_LOCK_TIMEOUT_CONTENTION_STDERR] * 6 + ["ok"], + ) + + assert harness.run() is True + assert len(harness.deploy_calls) == 7 + assert harness.rolled_back == [] + assert harness.sleeps == [] + + def test_a_p3009_row_logging_a_lock_timeout_rolls_back_and_retries(self, monkeypatch, tmp_path): + locked_row = _LedgerRow( + _LOCK_TIMEOUT_MIGRATION, + _P3009_STARTED_AT, + logs=_P3018_LOCK_TIMEOUT_STDERR, + ) + harness = _MigrateDeployHarness( + monkeypatch, + tmp_path, + [_p3009_stderr(_LOCK_TIMEOUT_MIGRATION, _P3009_STARTED_AT), "ok"], + ledger=_FakeLedger(at_error=(locked_row,), after_peer=(locked_row,)), + ) + + assert harness.run() is True + assert len(harness.deploy_calls) == 2 + assert harness.rolled_back == [_LOCK_TIMEOUT_MIGRATION] + + def test_a_lock_timeout_with_a_permission_error_stays_fatal(self, monkeypatch, tmp_path): + stderr = _P3018_LOCK_TIMEOUT_STDERR + "ERROR: permission denied for table LiteLLM_SpendLogs\n" + harness = _MigrateDeployHarness(monkeypatch, tmp_path, [stderr]) + + with pytest.raises(RuntimeError, match="insufficient"): + harness.run() + assert len(harness.deploy_calls) == 1 + assert harness.rolled_back == [] + + +class TestLogMigrationLockHolders: + _MIGRATION_SQL = ( + "-- AlterTable\n" + 'ALTER TABLE "LiteLLM_SpendLogs" ADD COLUMN IF NOT EXISTS "created_at" TIMESTAMP(3);\n' + ) + + def _migration_file(self, tmp_path: Path) -> None: + migration_dir = tmp_path / "migrations" / _LOCK_TIMEOUT_MIGRATION + migration_dir.mkdir(parents=True) + (migration_dir / "migration.sql").write_text(self._MIGRATION_SQL) + + def test_the_current_lock_holders_are_logged(self, monkeypatch, tmp_path, caplog): + self._migration_file(tmp_path) + monkeypatch.setenv("LITELLM_MIGRATION_DIR", str(tmp_path)) + monkeypatch.setenv("DATABASE_URL", "postgresql://u:p@localhost:9/x") + executed: list[tuple[object, object]] = [] + connect_kwargs: list[dict[str, object]] = [] + + class _Conn: + def __enter__(self): + return self + + def __exit__(self, *args): + return None + + def execute(self, query, params): + executed.append((query, params)) + row = ( + 1234, + "LiteLLM_SpendLogs", + "AccessExclusiveLock", + "idle in transaction", + "proxy-1", + "0:00:42", + ) + return type("Cur", (), {"fetchall": staticmethod(lambda: (row,))})() + + def fake_connect(*args, **kwargs): + connect_kwargs.append(kwargs) + return _Conn() + + monkeypatch.setattr("psycopg.connect", fake_connect) + + with caplog.at_level("WARNING", logger="litellm_proxy_extras"): + ProxyExtrasDBManager._log_migration_lock_holders(_LOCK_TIMEOUT_MIGRATION) + + assert len(executed) == 1 + assert connect_kwargs[0]["options"] == "-c statement_timeout=5000" + assert executed[0][1] == ("public", ["LiteLLM_SpendLogs"]) + assert "pg_locks" in str(executed[0][0]) + assert "a.query" not in str(executed[0][0]) + assert "1234" in caplog.text + assert "LiteLLM_SpendLogs" in caplog.text + assert "AccessExclusiveLock" in caplog.text + assert "0:00:42" in caplog.text + + def test_a_database_that_cannot_be_reached_is_skipped_quietly(self, monkeypatch, tmp_path): + import psycopg + + self._migration_file(tmp_path) + monkeypatch.setenv("LITELLM_MIGRATION_DIR", str(tmp_path)) + monkeypatch.setenv("DATABASE_URL", "postgresql://u:p@localhost:9/x") + + def refuse(*args, **kwargs): + raise psycopg.OperationalError("connection refused") + + monkeypatch.setattr("psycopg.connect", refuse) + + assert ProxyExtrasDBManager._log_migration_lock_holders(_LOCK_TIMEOUT_MIGRATION) is None + + def test_a_migration_without_sql_on_disk_never_connects(self, monkeypatch, tmp_path): + monkeypatch.setenv("LITELLM_MIGRATION_DIR", str(tmp_path)) + monkeypatch.setenv("DATABASE_URL", "postgresql://u:p@localhost:9/x") + connections = [] + monkeypatch.setattr("psycopg.connect", lambda *args, **kwargs: connections.append(1)) + + ProxyExtrasDBManager._log_migration_lock_holders("99999999999999_not_on_disk") + + assert connections == [] + + @pytest.mark.parametrize( "steps,logs,script,expected", (