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__) engine = create_engine( settings.database_url, pool_pre_ping=True, future=True, # #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