mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-02 02:11:58 +00:00
fix(proxy-extras): bound the lock waits of the partitioned SpendLogs index build (#44109)
Co-authored-by: yassin <yassin@berri.ai> Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
3f39fef52f
commit
123237c8a8
2 changed files with 185 additions and 21 deletions
|
|
@ -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}")
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue