"""#3463 — у запросов к БД обязан быть потолок по времени и по ожиданию блокировки. Отказ, который тут закрывается (см. issue): под `ACCESS EXCLUSIVE` на таблице шаг БД на пути `/estimate` ЖДЁТ блокировку; бюджет источника истекает, обёртка `run_db_thread` уходит ждать свой поток (иначе он останется сиротой в общей `Session` — #3449), а верхняя граница этого ожидания = длительность самого запроса. Границы у запроса не было → слот `_estimate_slots` не возвращался → `_ESTIMATE_CONCURRENCY = 4` исчерпывался и `/estimate` отдавал 429 всем остальным. Потолок поэтому стоит на КОННЕКТЕ (libpq `options`), а не в питоновской обёртке: таймаут в обёртке вернул бы ровно ту сироту, ради которой писался #3449. Проверки по значению, а не по тексту: * потолки реально доехали до параметров подключения движка и согласованы с объявленными бюджетами `/estimate` (без БД — падают от снятия `connect_args`); * живая сессия этого движка сообщает оба потолка (без БД пропускается); * запрос длиннее потолка ОБРЫВАЕТСЯ за отведённое время, а не висит; * ожидание блокировки длиннее потолка обрывается — это и есть сценарий #3463; * задача со своим `SET LOCAL statement_timeout` (планировщик: 900 с в `app/tasks/listing_source_snapshot.py`) новым сессионным потолком НЕ обрезается. """ from __future__ import annotations import inspect import os import time from collections.abc import Iterator import psycopg import pytest from sqlalchemy import Engine, create_engine, text from sqlalchemy.exc import OperationalError from app.core.db import _LOCK_TIMEOUT_MS, _STATEMENT_TIMEOUT_MS, DB_CONNECT_ARGS, engine # Потолки для ПОВЕДЕНЧЕСКИХ проверок — намеренно маленькие: проверяется механизм # (потолок из `options` реально обрывает запрос / ожидание лока), а прод-ВЕЛИЧИНЫ # проверяет `test_live_session_reports_both_ceilings` на самом прод-движке. Иначе # каждый прогон сьюта стоил бы 30 с ожидания pg_sleep. _PROBE_STATEMENT_TIMEOUT_MS = 1_000 _PROBE_LOCK_TIMEOUT_MS = 500 _LOCK_PROBE_TABLE = "t3463_lock_probe" def _engine_connect_options(eng: Engine) -> str: """Строка libpq `options`, с которой движок РЕАЛЬНО открывает коннекты. `connect_args` в движке не хранятся полем: `create_engine` вливает их в `cparams` замыкания `pool._creator`. Читаем оттуда, а не из `DB_CONNECT_ARGS`, — иначе тест остался бы зелёным после снятия `connect_args=` у `create_engine`. Наружу отдаём ТОЛЬКО `options`: в `cparams` лежит пароль роли, и текст провалившегося assert'а уехал бы с ним в лог CI. """ creator = getattr(eng.pool, "_creator", None) assert creator is not None, "у пула движка нет _creator — SQLAlchemy сменила устройство" cparams = inspect.getclosurevars(creator).nonlocals.get("cparams") assert cparams is not None, ( "в замыкании pool._creator нет cparams — SQLAlchemy сменила устройство, " "проверку параметров подключения надо переписать, а не удалять" ) return str(cparams.get("options", "")) def _live_engine(connect_args: dict[str, str]) -> Engine | None: """Движок против живой Postgres с заданными `connect_args`, иначе None. Тот же способ добыть DSN, что у `_live_session()` в tests/test_house_dedup_merge.py и tests/test_purge_expired_trade_in_data.py: в CI Postgres есть (ci-tradein.yml), на ноутбуке без БД тест пропускается (учтён в tests/skip_allowlist.txt). """ dsn = os.environ.get("TEST_DATABASE_URL") or os.environ.get("DATABASE_URL", "") if not dsn or "localhost:5432/test" in dsn: return None # Сначала проба БЕЗ connect_args: она отделяет «сервера нет» (честный пропуск) # от «сервер есть, но наши `options` он не принял». Глушить второе нельзя — # именно так испорченное значение (`statement_timeout=30000zz`) проходило # зелёным: коннект падал, тест пропускался, запись в allowlist гасила сигнал, # а на проде это FATAL на КАЖДОМ коннекте. try: probe = create_engine(dsn, future=True) except Exception: return None try: with probe.connect() as conn: conn.execute(text("SELECT 1")) except Exception: return None finally: probe.dispose() eng = create_engine(dsn, future=True, connect_args=connect_args) with eng.connect() as conn: # НЕ под except: сервер живой, виноваты connect_args conn.execute(text("SELECT 1")) return eng @pytest.fixture def probe_engine() -> Iterator[Engine]: """Живой движок с МАЛЕНЬКИМИ потолками — проверка механизма, не прод-величин.""" eng = _live_engine( { "options": ( f"-c statement_timeout={_PROBE_STATEMENT_TIMEOUT_MS} " f"-c lock_timeout={_PROBE_LOCK_TIMEOUT_MS}" ) } ) if eng is None: pytest.skip("живой Postgres недоступен (DATABASE_URL-заглушка)") try: yield eng finally: eng.dispose() # ── Проводка и согласованность величин (без БД) ─────────────────────────────── def test_engine_opens_connections_with_both_ceilings() -> None: """Оба потолка доехали до параметров подключения ПРОД-движка. Фальсификация: убрать `connect_args=DB_CONNECT_ARGS` из `create_engine` — `options` станет пустой, тест краснеет. """ expected = f"-c statement_timeout={_STATEMENT_TIMEOUT_MS} -c lock_timeout={_LOCK_TIMEOUT_MS}" options = _engine_connect_options(engine) # РАВЕНСТВО, а не `in`: подстрочная проверка пропускала испорченный хвост # (`…=30000zz` содержит `…=30000`), а Postgres на такое значение отвечает # `FATAL: invalid value for parameter "statement_timeout"` — ни одного коннекта # ни в одном из трёх сервисов образа, полный отказ продукта. assert options == expected, ( f"движок открывает коннекты с options={options!r}, ожидалось {expected!r}: " "либо потолка нет вовсе (заблокированный запрос снова висит и жжёт слот " "/estimate, #3463), либо значение испорчено — тогда libpq отвергнет КАЖДЫЙ коннект" ) assert DB_CONNECT_ARGS["options"] == options def test_ceilings_are_coherent_with_declared_estimate_budgets() -> None: """Потолок выше САМОГО ДЛИННОГО объявленного бюджета `/estimate`, но конечен.""" from app.core.config import settings declared_budgets_s = ( settings.estimate_avito_imv_timeout_s, # 20 с — самый длинный settings.estimate_geocode_budget_s, # 12 с settings.estimate_house_meta_timeout_s, # 8 с settings.estimate_yandex_valuation_timeout_s, # 8 с settings.estimate_cian_valuation_timeout_s, # 8 с ) longest_ms = max(declared_budgets_s) * 1000 assert _STATEMENT_TIMEOUT_MS > longest_ms, ( f"потолок запроса {_STATEMENT_TIMEOUT_MS} мс не выше самого длинного объявленного " f"бюджета {longest_ms:.0f} мс — потолок стал бы биндящим ограничением честной работы" ) # Потолок сверху — иначе у проверки есть пол и нет крыши: `_STATEMENT_TIMEOUT_MS = # 300_000` (пять минут) зеленел бы, а пять минут ожидания это тот же отказ, только # медленнее: четыре таких запроса всё так же выедают `_ESTIMATE_CONCURRENCY`. assert _STATEMENT_TIMEOUT_MS <= 2 * longest_ms, ( f"потолок запроса {_STATEMENT_TIMEOUT_MS} мс больше чем вдвое превышает самый " f"длинный объявленный бюджет {longest_ms:.0f} мс — это уже не защита, а отсрочка: " "слот /estimate держится всё это время" ) assert 0 < _LOCK_TIMEOUT_MS < _STATEMENT_TIMEOUT_MS, ( "ожидание блокировки обязано обрываться РАНЬШЕ потолка на сам запрос: " "деградировать лучше, чем держать слот" ) # ── Поведение против живой Postgres ────────────────────────────────────────── def test_live_session_reports_both_ceilings() -> None: """Прод-ВЕЛИЧИНЫ на живой сессии ПРОД-движка: сервер их принял, а не проигнорировал. Коннект берётся у самого `app.core.db.engine`, а не у собранного здесь двойника: двойник остался бы зелёным после снятия `connect_args=` в `create_engine`. """ if _live_engine(DB_CONNECT_ARGS) is None: pytest.skip("живой Postgres недоступен (DATABASE_URL-заглушка)") with engine.connect() as conn: statement_timeout = conn.execute( text("SELECT current_setting('statement_timeout')") ).scalar_one() lock_timeout = conn.execute(text("SELECT current_setting('lock_timeout')")).scalar_one() assert statement_timeout == "30s", ( f"сессия сообщает statement_timeout={statement_timeout!r} — " f"ожидалось 30s (_STATEMENT_TIMEOUT_MS={_STATEMENT_TIMEOUT_MS})" ) assert lock_timeout == "5s", ( f"сессия сообщает lock_timeout={lock_timeout!r} — " f"ожидалось 5s (_LOCK_TIMEOUT_MS={_LOCK_TIMEOUT_MS})" ) def test_statement_over_ceiling_is_cancelled_not_hung(probe_engine: Engine) -> None: """`pg_sleep` длиннее потолка обрывается ОТМЕНОЙ за отведённое время.""" sleep_s = _PROBE_STATEMENT_TIMEOUT_MS / 1000 * 5 started = time.monotonic() with probe_engine.connect() as conn, pytest.raises(OperationalError) as excinfo: conn.execute(text(f"SELECT pg_sleep({sleep_s})")) elapsed = time.monotonic() - started assert isinstance(excinfo.value.orig, psycopg.errors.QueryCanceled), ( f"запрос упал не отменой по таймауту, а {type(excinfo.value.orig).__name__}" ) assert elapsed < sleep_s, ( f"запрос шёл {elapsed:.1f} с при потолке {_PROBE_STATEMENT_TIMEOUT_MS} мс — " "потолок не сработал, он висел до конца pg_sleep" ) def test_lock_wait_over_ceiling_is_aborted(probe_engine: Engine) -> None: """Сценарий #3463: под ACCESS EXCLUSIVE читатель ОТВАЛИВАЕТСЯ, а не ждёт вечно.""" with probe_engine.connect() as blocker: blocker.execute(text(f"CREATE TABLE IF NOT EXISTS {_LOCK_PROBE_TABLE} (id int)")) blocker.commit() try: blocker.execute(text(f"LOCK TABLE {_LOCK_PROBE_TABLE} IN ACCESS EXCLUSIVE MODE")) started = time.monotonic() with probe_engine.connect() as victim, pytest.raises(OperationalError) as excinfo: victim.execute(text(f"SELECT count(*) FROM {_LOCK_PROBE_TABLE}")) elapsed = time.monotonic() - started finally: blocker.rollback() blocker.execute(text(f"DROP TABLE IF EXISTS {_LOCK_PROBE_TABLE}")) blocker.commit() assert isinstance(excinfo.value.orig, psycopg.errors.LockNotAvailable), ( f"читатель упал не по ожиданию блокировки, а {type(excinfo.value.orig).__name__} — " "сработал не тот потолок" ) assert elapsed < _PROBE_STATEMENT_TIMEOUT_MS / 1000, ( f"ожидание блокировки длилось {elapsed:.2f} с при lock_timeout " f"{_PROBE_LOCK_TIMEOUT_MS} мс — оборвал не lock_timeout" ) def test_set_local_statement_timeout_overrides_session_ceiling(probe_engine: Engine) -> None: """Задача со своим `SET LOCAL` НЕ обрезается сессионным потолком. Это и есть проверка обещания «задачи планировщика не пострадают»: у `app/tasks/listing_source_snapshot.py:288` стоит `SET LOCAL statement_timeout = 900000`, и он обязан ПЕРЕКРЫВАТЬ значение из `connect_args`, а не наоборот. """ own_budget_ms = _PROBE_STATEMENT_TIMEOUT_MS * 10 sleep_s = _PROBE_STATEMENT_TIMEOUT_MS / 1000 * 2 # заведомо больше сессионного потолка with probe_engine.connect() as conn: conn.execute(text(f"SET LOCAL statement_timeout = {own_budget_ms}")) effective = conn.execute(text("SELECT current_setting('statement_timeout')")).scalar_one() assert effective == "10s", f"SET LOCAL не применился: current_setting={effective!r}" # Не просто current_setting: запрос длиннее СЕССИОННОГО потолка обязан дойти до конца. conn.execute(text(f"SELECT pg_sleep({sleep_s})")) conn.rollback() def test_set_local_is_scoped_to_its_transaction(probe_engine: Engine) -> None: """Обратная сторона: чужой `SET LOCAL` не снимает потолок со всей сессии. Иначе одна задача с 900-секундным бюджетом отключала бы защиту у всех, кому достанется тот же коннект из пула. """ with probe_engine.connect() as conn: conn.execute(text(f"SET LOCAL statement_timeout = {_PROBE_STATEMENT_TIMEOUT_MS * 10}")) conn.rollback() after = conn.execute(text("SELECT current_setting('statement_timeout')")).scalar_one() assert after == "1s", ( f"после завершения транзакции statement_timeout={after!r} — " "SET LOCAL протёк за пределы своей транзакции, коннект вернулся в пул без потолка" )