import asyncio import logging from collections.abc import Callable, Generator from typing import Any from sqlalchemy import create_engine from sqlalchemy.orm import DeclarativeBase, Session, sessionmaker from app.core.config import settings logger = logging.getLogger(__name__) # #3463. Потолок ОДНОГО statement'а, секунды·1000. Ставится на КОННЕКТЕ (libpq # `options`), а не в питоновской обёртке: обёртка (`run_db_thread` ниже) при # отмене обязана ДОЖДАТЬСЯ потока, иначе поток остаётся сиротой в общей # `Session` — ровно то, ради чего писался #3449. Значит верхняя граница ожидания # = длительность самого запроса, и задать её может только сервер. # # 30 с выбраны так, чтобы потолок НИКОГДА не стал биндящим ограничением для # честной работы, но остался конечным: # * самый длинный ОБЪЯВЛЕННЫЙ бюджет на `/estimate` — 20 с (`estimate_avito_imv_timeout_s`, # config.py:852); дальше 12 с геокод, 8 с Yandex/Cian/house_meta. 30 с = 1.5× от максимума; # * ОСНОВНАЯ опора по планировщику — `scrape_runs` (длительности целых прогонов, они # не вытесняются): самая долгая ЧИСТО-БД задача за 14 суток — listing_source_snapshot, # 9.7 с ЦЕЛИКОМ (и у неё сверх того свой `SET LOCAL statement_timeout = 900000`, # который перекрывает это значение — гейт tests/test_3463_db_timeouts.py); # * самый длинный set-based statement ЧЕРЕЗ движок из замеренных — матч ГАР→houses # (`services/gar_flats_loader._MATCH_SQL`): 2.07 с с городским фильтром и 6.46 с без # него (`city_filter=None`, флаг CLI). Запас ~3×, и это СЧИТАЮЩИЙ запрос, а не ждущий. # # `pg_stat_statements` опорой по планировщику НЕ является: при `max = 5000` он вытесняет # редкие записи (проверено 12.09 — `dealloc` вырос на единицу за десять минут, и из топа # пропали ВСЕ записи с `calls = 1`, включая `REFRESH MATERIALIZED VIEW` 30.85 с и KNN # `cadastral_geo_match` 2.45 с). Суточная задача до следующих суток там не доживает, так # что «самый долгий запрос 4.27 с» верно только для ВЫСОКОЧАСТОТНЫХ запросов. _STATEMENT_TIMEOUT_MS = 30_000 # Ожидание БЛОКИРОВКИ — заведомо меньше: ждать лок дольше секунд смысла нет, лучше # деградировать. 5 с — та же величина, что у миграций проекта # (`SET LOCAL lock_timeout = '5s'` в data/sql/250,251,260,272,277…), снизу ограничена # deadlock_timeout (на проде 1 с — сверено 12.09). Именно этот потолок закрывает # сценарий #3463: под `ACCESS EXCLUSIVE` на `geocode_cache` запрос ЖДЁТ лок, а не # считает, — statement_timeout тут только страховка от «считает вечно». _LOCK_TIMEOUT_MS = 5_000 # idle_in_transaction_session_timeout НАМЕРЕННО не трогаем: тем же движком живёт tgbot, # и `services/tgbot/bridge.py` держит транзакцию открытой ПОВЕРХ long-poll Telegram # (замер на проде 12.09, 3 пробы с шагом 7 с: одна и та же сессия, запрос # `SELECT value FROM tg_support_state …`, возраст транзакции циклически растёт до ~29 с). # Сессионный потолок на простой в транзакции ронял бы long-poll КАЖДЫЙ цикл — # гарантированно, а не в редком случае. DB_CONNECT_ARGS = { "options": f"-c statement_timeout={_STATEMENT_TIMEOUT_MS} -c lock_timeout={_LOCK_TIMEOUT_MS}" } engine = create_engine( settings.database_url, pool_pre_ping=True, future=True, # #3463. Накрывает ВСЕ три сервиса образа (backend / scraper / tgbot — один и тот # же `app.core.db`, см. docker-compose.prod.yml) и обе стороны: продуктовый путь # `/estimate` и задачи планировщика. Миграции идут мимо (psql из # .forgejo/workflows/deploy-tradein.yml, не этот движок) — их DDL под своим # `SET LOCAL lock_timeout` и потолком не ограничен. connect_args=DB_CONNECT_ARGS, # #3194: SQLAlchemy печатает ВСЕ bind-параметры в тексте StatementError — # через них в GlitchTip уезжали ключ шифрования кук и сами куки # (pgp_sym_encrypt(:cookies_json, :key)). Флаг на УРОВНЕ ДВИЖКА кроет все # сайты вызова разом, включая будущие. # НЕ закрывает: текст ошибки самого драйвера (Postgres DETAIL со значением) # и сырые psycopg-подключения мимо движков — это отдельный класс. hide_parameters=True, # #3408 п.1. Потолок пула обязан быть НЕ МЕНЬШЕ суммы потолков одновременности, # которые сам же процесс и объявляет: 4 оценки (`api/v1/trade_in` # `_ESTIMATE_CONCURRENCY`) + 4 подсказки (`api/public/mera` `_SUGGEST_CONCURRENCY`) # + 8 фоновых догрузок (`services/estimator` `_MAX_DEFERRED_REFRESH_TASKS`) = 16. # Дефолт SQLAlchemy 5+10=15 меньше этой суммы, то есть исчерпать пул можно # штатной работой, не абузом. Гейт — tests/test_3408_pool_ceiling.py. # # Инвариант «пул >= суммы объявленных потолков» — это ПОЛ, а не гарантия: # сумма считает по одному коннекту на объявленную единицу работы, а коннект # держит и всё остальное — любая ручка с `Depends(get_db)` тоже. Глобального # обработчика `sqlalchemy.exc.TimeoutError` в app/main.py нет, поэтому на # исчерпанном пуле соседние ручки отдают 500 (после ожидания `pool_timeout`), # а не медленный 200. # # pool_size оставлен дефолтным (5): это ПОСТОЯННО открытые коннекты, а в покое # прод держит 5-6 (замер 11.09). Растёт только overflow — коннекты пика, # которые пул закрывает сам. Потолок процесса: 5 + 15 = 20; воркер один # (docker-compose.prod.yml, uvicorn без --workers), Postgres max_connections=100. max_overflow=15, # #3408 п.2. Дефолтные 30 с ожидания коннекта длиннее ЛЮБОГО бюджета внешнего # источника в эстиматоре: 8 с (Yandex / Cian / house_meta), 12 с (geocode), # 20 с (Avito IMV — `estimate_avito_imv_timeout_s`, config.py:852). Сам чекаут # прервать `asyncio.wait_for` не может: он занимает поток `asyncio.to_thread` # целиком, а пул потоков конечен (min(32, cpu+4)) — исчерпанный пул коннектов # так превращается в исчерпанный пул потоков. 5 с короче самого КОРОТКОГО # бюджета: занятый пул деградирует ОДИН источник, а не весь запрос. # # Отдельный коммит в конце ветки намеренно (ревью PR #3444): это единственная # правка, которая меняет режим отказа с «медленно» на «быстро с ошибкой», и # едет она во ВСЕ сервисы образа — backend, scraper, tgbot # (tradein-mvp/docker-compose.prod.yml). За 29 ч логов исчерпания пула не было # ни разу, то есть новое значение на проде пока не на чем проверить. # ТРИГГЕР ОТКАТА на 30 с: любое `QueuePool limit ... timed out` в логах # бэкенда ЛИБО рост failed+zombie в `scrape_runs` после деплоя. pool_timeout=5, ) SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine, expire_on_commit=False) class Base(DeclarativeBase): pass def get_db() -> Generator[Session, None, None]: db = SessionLocal() try: yield db finally: db.close() async def run_db_thread[T](fn: Callable[..., T], *args: Any, **kwargs: Any) -> T: """`asyncio.to_thread(fn, ...)`, который при отмене ДОЖИДАЕТСЯ своего потока (#3449). Для синхронной работы по сессии, которая шагу НЕ принадлежит — она общая со всем остальным запросом (`Depends(get_db)`). Поток отменить нельзя: у `asyncio.to_thread` отменяется только ожидание со стороны loop'а. Корутина умирает по бюджету источника (`estimator._with_budget` = `asyncio.wait_for`, геокодер — 12 с), а поток продолжает работать с ТОЙ ЖЕ `Session`, пока вызывающий уже идёт дальше по коду — следующий источник, `_fetch_anchor_comps`, `_persist_estimate_and_commit`. Два потока в одной `Session` дают «another operation is in progress» / `InvalidRequestError` на СЛЕДУЮЩЕМ шаге: у источников такую ошибку глушит `except` вокруг вызова, у персиста оценки не глушит ничего — 500 и потерянная оценка клиента. Поэтому отмена пробрасывается ПОСЛЕ того, как поток отпустил сессию. Цена — бюджет источника переезжает на длину ОДНОГО шага БД (чекаут ≤ `pool_timeout` плюс сам запрос), а не на длину фетча, ради которой бюджет и заведён. Только защита от сироты: транзакцию шаг НЕ завершает (в середине геокодинга commit зафиксировал бы частичное состояние оценки). Кому нужен ещё и возврат коннекта в пул перед внешним HTTP — `estimator._db_step`, он поверх этого. Гейт — tests/test_3449_geocoder_cancel_orphan.py. """ step = asyncio.ensure_future(asyncio.to_thread(fn, *args, **kwargs)) try: return await asyncio.shield(step) except asyncio.CancelledError: await asyncio.wait([step]) if not step.cancelled() and step.exception() is not None: # Результата уже никто не ждёт: без явного чтения asyncio напечатает # «Task exception was never retrieved» вообще без контекста. logger.warning("шаг БД упал уже после отмены: %s", step.exception()) raise