diff --git a/tradein-mvp/backend/app/core/db.py b/tradein-mvp/backend/app/core/db.py index dd260dce..408d47ba 100644 --- a/tradein-mvp/backend/app/core/db.py +++ b/tradein-mvp/backend/app/core/db.py @@ -23,18 +23,18 @@ engine = create_engine( # Дефолт 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, - # Дефолтные 30 с ожидания коннекта — вчетверо больше любого бюджета внешнего - # источника в эстиматоре (8 с, `estimate_*_timeout_s`). Такой чекаут нельзя - # прервать `asyncio.wait_for`: он занимает поток `asyncio.to_thread` целиком, а - # пул потоков сам конечен (min(32, cpu+4)) — исчерпанный пул коннектов так - # превращается в исчерпанный пул потоков. 5 с < бюджета источника: занятый пул - # деградирует ОДИН источник, а не весь запрос. - pool_timeout=5, ) SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine, expire_on_commit=False) diff --git a/tradein-mvp/backend/app/services/estimator.py b/tradein-mvp/backend/app/services/estimator.py index 75bdfbc9..b757e7cd 100644 --- a/tradein-mvp/backend/app/services/estimator.py +++ b/tradein-mvp/backend/app/services/estimator.py @@ -696,9 +696,9 @@ YANDEX_VALUATION_DEFAULT_TYPE = "SELL" async def _db_step[T](db: Session, work: Callable[[], T]) -> T: - """Синхронный шаг БД внутри async-функции: в потоке И с возвратом коннекта в пул. + """Синхронный шаг БД внутри async-функции: в потоке, с возвратом коннекта и без сирот. - Две вещи сразу, потому что обе про одно и то же — `Session` из `get_db` + Три вещи сразу, потому что все про одно и то же — `Session` из `get_db` (#3408 п.1/п.2): 1. `db.execute()` — блокирующий вызов. На loop'е он держит ВЕСЬ инстанс @@ -710,6 +710,8 @@ async def _db_step[T](db: Session, work: Callable[[], T]) -> T: взятый ради одного SELECT'а кэша, живёт до конца функции — включая ожидание чужого HTTP (замер 11.09: 8.5 с на догрузку Яндекса), и пиковый спрос считается не в коротких SELECT'ах, а в целых фетчах. + 3. Отмена шага ДОЖИДАЕТСЯ потока (`asyncio.shield` ниже). Поток отменить + нельзя, а сессия у шага не своя — общая со всем остальным запросом. ЧТО ЭТО МЕНЯЕТ ДЛЯ ВЫЗЫВАЮЩЕГО: незакоммиченная работа его транзакции (например fill-null backfill ФИАСа в `estimate_quality`) фиксируется здесь, @@ -729,7 +731,30 @@ async def _db_step[T](db: Session, work: Callable[[], T]) -> T: db.commit() return result - return await asyncio.to_thread(_run) + step = asyncio.ensure_future(asyncio.to_thread(_run)) + try: + return await asyncio.shield(step) + except asyncio.CancelledError: + # Отмена по бюджету источника (`_with_budget` = `asyncio.wait_for`) поток НЕ + # останавливает: у `to_thread` отменяется только ожидание со стороны loop'а. + # Корутина умирает, поток продолжает работать с ТОЙ ЖЕ `Session`, а вызывающий + # тем временем идёт дальше — следующий источник, `_fetch_anchor_comps`, + # `_persist_estimate_and_commit` — ПО ТОЙ ЖЕ сессии. Два потока в одной сессии + # дают «another operation is in progress» / InvalidRequestError на следующем + # шаге БД: у источников такую ошибку глушит `except` вокруг вызова, у персиста + # оценки не глушит ничего — это 500 и потерянная оценка клиента, причём ровно + # под нагрузкой, ради которой правка и делается. + # + # Поэтому отмену пробрасываем ПОСЛЕ того, как поток отпустил сессию. Цена — + # бюджет источника переезжает на длину ОДНОГО шага БД (чекаут ≤ `pool_timeout` + # плюс сам запрос), а не на длину фетча, ради которой бюджет и заведён. + # Гейт — tests/test_3408_db_step_cancel_orphan.py. + await asyncio.wait([step]) + if not step.cancelled() and step.exception() is not None: + # Результата уже никто не ждёт: без явного чтения asyncio напечатает + # «Task exception was never retrieved» вообще без контекста. + logger.warning("db-шаг упал уже после отмены по бюджету: %s", step.exception()) + raise async def _get_or_fetch_imv_cached( @@ -940,9 +965,18 @@ _DEFERRED_REFRESH_TASKS: set[asyncio.Task[None]] = set() # 100 у Postgres. Догрузка почти всё время ждёт чужой HTTP, поэтому небольшой # потолок не занижает пропускную способность прогрева. # -# Коннект к БД задача держит только на самих SELECT/UPSERT кэша, а не весь фетч -# (`_db_step`, #3408): до правки одна догрузка Яндекса удерживала соединение -# 8.5 с (замер на проде 11.09), то есть восемь таких задач выедали пул целиком. +# Фоновых источника ДВА, и по удержанию коннекта они РАЗНЫЕ (#3408): +# • yandex_valuation — коннект только на самих SELECT/UPSERT кэша (`_db_step`): +# до правки одна догрузка держала соединение 8.5 с (замер на проде 11.09), +# то есть восемь таких задач выедали пул целиком; +# • cian_valuation — коннект на ВЕСЬ фетч (HTTP-таймаут 25 с): транзакцию +# открывают `_load_from_cache` / `load_session` на входе, а закрывает её +# `db.commit()` в самом конце (scraper_kit providers/cian/valuation.py:163, +# 176, 595). Известный остаток: функция общая со скраппером +# (`app/tasks/cian_history_backfill.py:458`), где коммит в середине менял бы +# семантику батча, — нужен отдельный опт-ин путь, а не правка на месте. +# На ПОТОЛОК пула остаток не влияет (задача держит один коннект независимо от +# того, как долго) — только на среднюю занятость пула. _MAX_DEFERRED_REFRESH_TASKS = 8 @@ -1157,8 +1191,9 @@ async def _get_or_fetch_yandex_valuation_cached( len(result.history_items), ) except Exception as e: + # Откат уже сделан в потоке (`_db_step` → `except` → `db.rollback()`), а здесь + # это был бы блокирующий вызов на event loop'е — ровно то, что правка убирает. logger.warning("yandex_valuation: cache save failed (continuing): %s", e) - db.rollback() return result diff --git a/tradein-mvp/backend/tests/test_3408_db_step_cancel_orphan.py b/tradein-mvp/backend/tests/test_3408_db_step_cancel_orphan.py new file mode 100644 index 00000000..e58c449e --- /dev/null +++ b/tradein-mvp/backend/tests/test_3408_db_step_cancel_orphan.py @@ -0,0 +1,82 @@ +"""Отмена шага БД по бюджету не оставляет ОСИРОТЕВШИЙ поток в чужой сессии (#3408). + +`_with_budget` — это `asyncio.wait_for`, а `asyncio.to_thread` отменить нельзя: +ожидание со стороны loop'а снимается, поток продолжает работать. Сессия у шага не +своя — та же самая, с которой запрос идёт дальше: следующий источник, +`_fetch_anchor_comps`, `_persist_estimate_and_commit`. Два потока в одной `Session` +дают «another operation is in progress» / InvalidRequestError на СЛЕДУЮЩЕМ шаге, +и у персиста оценки этой ошибки не ловит никто — 500 и потерянная оценка клиента. + +Меряем значение, а не форму: «следующий шаг не вошёл в сессию, пока сирота не +закончил». Следующий шаг здесь — ГОЛЫЙ `asyncio.to_thread(db...)`, как +`_persist_estimate_and_commit` (estimator.py), а не ещё один `_db_step`: защита, +которая живёт только внутри `_db_step`, ровно того пострадавшего и не закрывает. +""" + +from __future__ import annotations + +import asyncio +import os +import threading +import time +from typing import Any + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +from app.services import estimator as est + +# Шаг БД заметно длиннее бюджета — окно, в котором сирота ещё работает. +_STEP_S = 0.3 +_BUDGET_S = 0.05 + + +class _ConcurrencyProbeSession: + """Session-дублёр, который считает ОДНОВРЕМЕННЫЕ входы в сессию.""" + + def __init__(self) -> None: + self._guard = threading.Lock() + self._inside = 0 + self.conflicts = 0 + self.executed = 0 + self.commits = 0 + + def execute(self, *_a: Any, **_kw: Any) -> Any: + with self._guard: + self._inside += 1 + self.executed += 1 + if self._inside > 1: + self.conflicts += 1 + try: + time.sleep(_STEP_S) + finally: + with self._guard: + self._inside -= 1 + return self + + def commit(self) -> None: + self.commits += 1 + + def rollback(self) -> None: # pragma: no cover — на этом пути не звучит + pass + + +async def test_budget_cancel_does_not_leave_orphan_in_session() -> None: + db = _ConcurrencyProbeSession() + + degraded = await est._with_budget( + est._db_step(db, lambda: db.execute("cache lookup")), # type: ignore[arg-type] + _BUDGET_S, + label="проба", + ) + # Бюджет истёк — источник деградировал в None, вызывающий идёт дальше. + assert degraded is None, "бюджет не сработал — тест ничего не проверил" + + # Следующий шаг ТОГО ЖЕ запроса по ТОЙ ЖЕ сессии (образец — persist оценки). + await asyncio.to_thread(db.execute, "persist estimate") + + assert db.executed == 2, f"звучали не оба шага (executed={db.executed})" + assert db.conflicts == 0, ( + "следующий шаг вошёл в сессию, пока осиротевший поток ещё работал в ней: " + "два потока в одной Session → «another operation is in progress» на персисте" + ) + assert db.commits >= 1, "сирота не завершил транзакцию — коннект не вернулся в пул" diff --git a/tradein-mvp/backend/tests/test_3408_pool_ceiling.py b/tradein-mvp/backend/tests/test_3408_pool_ceiling.py index 9cbde0ef..254340db 100644 --- a/tradein-mvp/backend/tests/test_3408_pool_ceiling.py +++ b/tradein-mvp/backend/tests/test_3408_pool_ceiling.py @@ -16,6 +16,7 @@ from __future__ import annotations import os +from typing import Any os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") @@ -25,30 +26,26 @@ from app.core.db import engine from app.services.estimator import _MAX_DEFERRED_REFRESH_TASKS +def _max_overflow(pool: Any) -> int: + """`max_overflow` у QueuePool публичного геттера не имеет (в отличие от + `size()` / `timeout()`), поэтому читаем приватное поле — но с явным падением: + молча переименуется в новой SQLAlchemy → гейт бы просто перестал что-либо + проверять (AttributeError упал бы ошибкой теста, а вот подмена значения на + дефолт — нет). + """ + assert hasattr(pool, "_max_overflow"), ( + f"{type(pool).__name__} больше не хранит `_max_overflow` — SQLAlchemy сменила " + "API пула, потолок коннектов читать больше нечем: почини гейт, а не удаляй его" + ) + return int(pool._max_overflow) + + def test_pool_ceiling_covers_declared_concurrency() -> None: declared = _ESTIMATE_CONCURRENCY + _SUGGEST_CONCURRENCY + _MAX_DEFERRED_REFRESH_TASKS - ceiling = engine.pool.size() + engine.pool._max_overflow + ceiling = engine.pool.size() + _max_overflow(engine.pool) assert ceiling >= declared, ( f"пул {ceiling} меньше суммы потолков одновременности {declared} " f"({_ESTIMATE_CONCURRENCY} оценок + {_SUGGEST_CONCURRENCY} подсказок + " f"{_MAX_DEFERRED_REFRESH_TASKS} фоновых догрузок) — штатная работа исчерпает пул" ) - - -def test_pool_checkout_wait_shorter_than_source_budget() -> None: - """Ожидание коннекта короче бюджета внешнего источника. - - Иначе занятый пул съедает весь бюджет запроса (и поток `asyncio.to_thread`, - которых тоже конечное число) вместо того, чтобы деградировать один источник. - """ - from app.core.config import settings - - budget = min( - settings.estimate_yandex_valuation_timeout_s, - settings.estimate_cian_valuation_timeout_s, - ) - - assert engine.pool._timeout < budget, ( - f"pool_timeout {engine.pool._timeout}с ≥ бюджета источника {budget}с" - )