mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-11 03:38:38 +00:00
fix(proxy-extras): bound migration DDL lock waits with lock_timeout and retry (#45389)
Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
b72bfa735b
commit
cf1284f3d5
4 changed files with 427 additions and 21 deletions
|
|
@ -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),
|
||||
}
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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."""
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
(
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue