From 123237c8a81936767387507f0ba4aff9fd48884e Mon Sep 17 00:00:00 2001 From: "devin-ai-integration[bot]" <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Thu, 1 Oct 2026 18:28:50 -0700 Subject: [PATCH] fix(proxy-extras): bound the lock waits of the partitioned SpendLogs index build (#44109) Co-authored-by: yassin Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../request_log_indexes.py | 57 ++++--- .../test_request_log_indexes.py | 149 ++++++++++++++++++ 2 files changed, 185 insertions(+), 21 deletions(-) diff --git a/litellm-proxy-extras/litellm_proxy_extras/request_log_indexes.py b/litellm-proxy-extras/litellm_proxy_extras/request_log_indexes.py index a344684a4cd..6c31e8364a9 100644 --- a/litellm-proxy-extras/litellm_proxy_extras/request_log_indexes.py +++ b/litellm-proxy-extras/litellm_proxy_extras/request_log_indexes.py @@ -56,8 +56,10 @@ REQUEST_LOG_INDEXES: Final = ( ) _IDENTIFIER_MAX_BYTES: Final = 63 -_PARENT_LOCK_TIMEOUT: Final = "2s" -_PARENT_LOCK_ATTEMPTS: Final = 30 +_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( @@ -181,6 +183,32 @@ def _under_migration_lock(connection: "psycopg.Connection[tuple[object, ...]]", 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 @@ -368,12 +396,13 @@ def build_index_on_partitioned_table( "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 _under_migration_lock( + 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) @@ -393,29 +422,15 @@ def _create_parent_index( table: str, index: RequestLogIndex, ) -> bool: - """Create the metadata-only parent index. Postgres takes a SHARE lock on the - parent for that statement, so it waits for in-flight writes and queues new ones - behind it; a short lock_timeout with retries keeps every such pause bounded.""" - import psycopg + """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(sql.SQL("SET lock_timeout = {}").format(sql.Literal(_PARENT_LOCK_TIMEOUT))) - try: - for _ in range(_PARENT_LOCK_ATTEMPTS): - try: - connection.execute(statement) - return True - except psycopg.errors.LockNotAvailable: - logger.info("Waiting for in-flight writes to %s before creating the parent index %s", table, name) - time.sleep(random.uniform(0.1, 0.5)) - finally: - connection.execute("SET lock_timeout = 0") - logger.warning("Could not get the parent lock on %s to create %s, leaving it for the next index build", table, name) - return False + connection.execute(statement) + return True def _attach_child_index( @@ -445,4 +460,4 @@ def _attach_child_index( logger.info("Attached index %s on partition %s to %s", child_index, child.name, parent_index) return True - return _under_migration_lock(connection, attach) + return _with_bounded_lock(connection, attach, f"attaching {child_index}") diff --git a/tests/proxy_migration_tests/test_request_log_indexes.py b/tests/proxy_migration_tests/test_request_log_indexes.py index 3b0c87b2a97..23e4adce477 100644 --- a/tests/proxy_migration_tests/test_request_log_indexes.py +++ b/tests/proxy_migration_tests/test_request_log_indexes.py @@ -14,6 +14,7 @@ from typing import Final import psycopg import pytest +from litellm_proxy_extras import request_log_indexes from litellm_proxy_extras.migration_lock import MIGRATION_LOCK_KEY, migration_lock from litellm_proxy_extras.migration_recovery import roll_back_failed_inert_migration from litellm_proxy_extras.request_log_indexes import ( @@ -649,6 +650,154 @@ def test_inserts_keep_flowing_while_the_partition_indexes_build(partitioned_data _assert_index_covers_every_partition(partitioned_database, CALL_ID_INDEX, "litellm_call_id_idx") +def _insert_for(database_url: str, seconds: float) -> None: + with psycopg.connect(database_url, autocommit=True) as conn: + conn.execute("SET lock_timeout = '1s'") + deadline: Final = time.monotonic() + seconds + while time.monotonic() < deadline: + _insert_spend_log(conn, f"lock-test-{uuid.uuid4().hex}", "2026-08-16") + time.sleep(0.05) + + +def _wait_for_blocked_ddl(database_url: str, query_pattern: str) -> bool: + with psycopg.connect(database_url, autocommit=True) as conn: + deadline: Final = time.monotonic() + 10 + while time.monotonic() < deadline: + if conn.execute( + "SELECT 1 FROM pg_stat_activity WHERE wait_event_type = 'Lock' AND query ILIKE %s", + (query_pattern,), + ).fetchone(): + return True + time.sleep(0.01) + return False + + +@requires_db +def test_inserts_are_never_held_back_while_the_parent_index_waits_for_an_open_write( + partitioned_database: str, +) -> None: + outcome: Final[list[bool]] = [] # mutable-ok: the builder thread hands its result back through it + with psycopg.connect(partitioned_database) as writer: + _insert_spend_log(writer, "parent-index-lock-owner", "2026-08-15") + builder_thread: Final = threading.Thread( + target=lambda: outcome.append(_build_in_its_own_session(partitioned_database, CALL_ID_INDEX_DEFINITION)) + ) + builder_thread.start() + try: + assert _wait_for_blocked_ddl(partitioned_database, "%CREATE INDEX%ON ONLY%") + _insert_for(partitioned_database, 3) + finally: + try: + writer.commit() + finally: + builder_thread.join() + assert outcome == [True] + _assert_index_covers_every_partition(partitioned_database, CALL_ID_INDEX, "litellm_call_id_idx") + + +@requires_db +def test_inserts_are_never_held_back_while_attach_partition_waits_for_a_reader_of_the_child_index( + partitioned_database: str, +) -> None: + partition: Final = "LiteLLM_SpendLogs_p2026_08" + child_index: Final = CALL_ID_INDEX_DEFINITION.partition_index_name(partition) + assert child_index == "LiteLLM_SpendLogs_p2026_08_litellm_call_id_idx" + with psycopg.connect(partitioned_database, autocommit=True) as conn: + conn.execute( + 'CREATE INDEX "LiteLLM_SpendLogs_litellm_call_id_idx" ON ONLY "LiteLLM_SpendLogs" ("litellm_call_id")' + ) + conn.execute( + sql.SQL('CREATE INDEX {} ON {} ("litellm_call_id")').format( + sql.Identifier(CALL_ID_INDEX_DEFINITION.partition_index_name(partition)), sql.Identifier(partition) + ) + ) + + outcome: Final[list[bool]] = [] # mutable-ok: the builder thread hands its result back through it + with psycopg.connect(partitioned_database) as reader: + reader.execute("SET enable_seqscan = off") + reader.execute( + sql.SQL('SELECT count(*) FROM {} WHERE "litellm_call_id" IS NULL').format(sql.Identifier(partition)) + ).fetchone() + reader_pid: Final = reader.execute("SELECT pg_backend_pid()").fetchone()[0] + with psycopg.connect(partitioned_database, autocommit=True) as inspector: + child_lock: Final = inspector.execute( + "SELECT 1 FROM pg_locks WHERE pid = %s AND relation = to_regclass(%s) " + "AND mode = 'AccessShareLock' AND granted", + (reader_pid, f'"{child_index}"'), + ).fetchone() + assert child_lock is not None + builder_thread: Final = threading.Thread( + target=lambda: outcome.append(_build_in_its_own_session(partitioned_database, CALL_ID_INDEX_DEFINITION)) + ) + builder_thread.start() + try: + assert _wait_for_blocked_ddl(partitioned_database, "%ATTACH PARTITION%") + _insert_for(partitioned_database, 3) + finally: + try: + reader.commit() + finally: + builder_thread.join() + assert outcome == [True] + _assert_index_covers_every_partition(partitioned_database, CALL_ID_INDEX, "litellm_call_id_idx") + + +@requires_db +def test_a_parent_index_that_never_gets_its_lock_is_left_for_the_next_index_build( + partitioned_database: str, monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture +) -> None: + monkeypatch.setattr(request_log_indexes, "_DDL_LOCK_ATTEMPTS", 2) + with psycopg.connect(partitioned_database) as writer: + _insert_spend_log(writer, "parent-index-lock-owner", "2026-08-15") + with caplog.at_level("WARNING", logger="litellm_proxy_extras"): + assert _build_in_its_own_session(partitioned_database, CALL_ID_INDEX_DEFINITION) is False + assert "leaving it for the next index build" in caplog.text + assert _indexed_table(partitioned_database, CALL_ID_INDEX) is None + writer.commit() + assert _build_in_its_own_session(partitioned_database, CALL_ID_INDEX_DEFINITION) is True + _assert_index_covers_every_partition(partitioned_database, CALL_ID_INDEX, "litellm_call_id_idx") + + +@requires_db +def test_an_attach_that_never_gets_its_lock_is_left_for_the_next_index_build( + partitioned_database: str, monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture +) -> None: + partition: Final = "LiteLLM_SpendLogs_p2026_08" + child_index: Final = CALL_ID_INDEX_DEFINITION.partition_index_name(partition) + assert child_index == "LiteLLM_SpendLogs_p2026_08_litellm_call_id_idx" + with psycopg.connect(partitioned_database, autocommit=True) as conn: + conn.execute( + 'CREATE INDEX "LiteLLM_SpendLogs_litellm_call_id_idx" ON ONLY "LiteLLM_SpendLogs" ("litellm_call_id")' + ) + conn.execute( + sql.SQL('CREATE INDEX {} ON {} ("litellm_call_id")').format( + sql.Identifier(child_index), sql.Identifier(partition) + ) + ) + + monkeypatch.setattr(request_log_indexes, "_DDL_LOCK_ATTEMPTS", 2) + with psycopg.connect(partitioned_database) as reader: + reader.execute("SET enable_seqscan = off") + reader.execute( + sql.SQL('SELECT count(*) FROM {} WHERE "litellm_call_id" IS NULL').format(sql.Identifier(partition)) + ).fetchone() + reader_pid: Final = reader.execute("SELECT pg_backend_pid()").fetchone()[0] + with psycopg.connect(partitioned_database, autocommit=True) as inspector: + child_lock: Final = inspector.execute( + "SELECT 1 FROM pg_locks WHERE pid = %s AND relation = to_regclass(%s) " + "AND mode = 'AccessShareLock' AND granted", + (reader_pid, f'"{child_index}"'), + ).fetchone() + assert child_lock is not None + with caplog.at_level("WARNING", logger="litellm_proxy_extras"): + assert _build_in_its_own_session(partitioned_database, CALL_ID_INDEX_DEFINITION) is False + assert "Could not get the lock for attaching" in caplog.text + reader.commit() + + assert _build_in_its_own_session(partitioned_database, CALL_ID_INDEX_DEFINITION) is True + _assert_index_covers_every_partition(partitioned_database, CALL_ID_INDEX, "litellm_call_id_idx") + + def _build_in_its_own_session(database_url: str, index: RequestLogIndex) -> bool: with psycopg.connect(database_url, autocommit=True) as builder: return build_index_on_partitioned_table(builder, "public", index)