From 9f696299de9f215984767254492201e777b3b5fe Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 12 Sep 2026 01:01:54 +0500 Subject: [PATCH 1/3] =?UTF-8?q?fix(mera):=20sync-=D0=91=D0=94=20=D0=B8?= =?UTF-8?q?=D1=81=D1=82=D0=BE=D1=87=D0=BD=D0=B8=D0=BA=D0=BE=D0=B2=20=D1=8D?= =?UTF-8?q?=D1=81=D1=82=D0=B8=D0=BC=D0=B0=D1=82=D0=BE=D1=80=D0=B0=20?= =?UTF-8?q?=E2=80=94=20=D1=81=20event=20loop=20=D0=B2=20=D0=BF=D0=BE=D1=82?= =?UTF-8?q?=D0=BE=D0=BA=20=D0=B8=20=D0=BD=D0=B5=20=D1=87=D0=B5=D1=80=D0=B5?= =?UTF-8?q?=D0=B7=20=D1=84=D0=B5=D1=82=D1=87?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Замер на проде 11.09 (изнутри хоста, тот же контейнер): - одна оценка 0.44 с (повтор адреса) / 0.97 с (новый адрес), из них БД 252/458 мс; - N=8 параллельных — все 200, heartbeat /health p95 5-7 мс, max 113-160 мс: loop сегодня НЕ голодает, «добавить воркеров uvicorn» замером не подтверждается (и умножило бы на N оба семафора, пять in-process лимитеров и пул); - зато одна фоновая догрузка Яндекса держала коннект пула 8.5 с (лиз прокси 33.856 → запись 42.334), а таких задач разрешено 8 при пуле 15. Правки: - estimator `_db_step`: SELECT/UPSERT кэша источников уходят в `asyncio.to_thread` и завершают транзакцию — коннект возвращается в пул ДО внешнего HTTP; - core/db: max_overflow 10→15 (потолок 20 на процесс ≥ 4+4+8 объявленных потолков одновременности) и pool_timeout 30→5 с (короче бюджета источника 8 с, иначе занятый пул съедает и бюджет запроса, и поток to_thread). Локальный замер ДО/ПОСЛЕ на тех же величинах: loop стоял 301 мс (0 тиков соседней корутины) → 0.2 мс (23.5k тиков); ожидание коннекта соседом во время фетча — таймаут пула → 0.1 мс. Refs #3083, #3408 --- tradein-mvp/backend/app/api/v1/trade_in.py | 6 +- tradein-mvp/backend/app/core/db.py | 19 ++ tradein-mvp/backend/app/services/estimator.py | 115 +++++++--- .../tests/test_3408_estimator_db_off_loop.py | 213 ++++++++++++++++++ .../backend/tests/test_3408_pool_ceiling.py | 54 +++++ 5 files changed, 374 insertions(+), 33 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_3408_estimator_db_off_loop.py create mode 100644 tradein-mvp/backend/tests/test_3408_pool_ceiling.py diff --git a/tradein-mvp/backend/app/api/v1/trade_in.py b/tradein-mvp/backend/app/api/v1/trade_in.py index bc05717c..512785ab 100644 --- a/tradein-mvp/backend/app/api/v1/trade_in.py +++ b/tradein-mvp/backend/app/api/v1/trade_in.py @@ -79,8 +79,10 @@ _estimate_limiter = SlidingWindowLimiter( # внешние тиры; пила одновременных оценок выедает пул и тормозит весь /api/v1/*. # Образец — public/mera.py::_suggest_slots (4 слота на секундное автодополнение). # -# 4 слота: вместе с 4 слотами suggest — 8 одновременно удерживаемых соединений -# из 15 возможных, остаток пула остаётся прочим ручкам. Ожидание слота 5с ≈ две +# 4 слота: вместе с 4 слотами suggest — 8 одновременно удерживаемых соединений. +# Пул под это заведомо шире: 5+15=20 на процесс, и потолок пула держится не +# меньше СУММЫ объявленных потолков одновременности, включая 8 фоновых догрузок +# (core/db.py + tests/test_3408_pool_ceiling.py, #3408). Ожидание слота 5с ≈ две # длительности оценки: если за это время слот не освободился, очередь глубока и # честный ответ — быстрый 429 с Retry-After, а не растущая очередь (очередь под # нагрузкой — те же занятые соединения плюс таймаут у клиента; mera.py:117-127). diff --git a/tradein-mvp/backend/app/core/db.py b/tradein-mvp/backend/app/core/db.py index cd18b717..dd260dce 100644 --- a/tradein-mvp/backend/app/core/db.py +++ b/tradein-mvp/backend/app/core/db.py @@ -16,6 +16,25 @@ engine = create_engine( # НЕ закрывает: текст ошибки самого драйвера (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. + # + # 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 12f0aa31..75bdfbc9 100644 --- a/tradein-mvp/backend/app/services/estimator.py +++ b/tradein-mvp/backend/app/services/estimator.py @@ -695,6 +695,43 @@ YANDEX_VALUATION_DEFAULT_CATEGORY = "APARTMENT" YANDEX_VALUATION_DEFAULT_TYPE = "SELL" +async def _db_step[T](db: Session, work: Callable[[], T]) -> T: + """Синхронный шаг БД внутри async-функции: в потоке И с возвратом коннекта в пул. + + Две вещи сразу, потому что обе про одно и то же — `Session` из `get_db` + (#3408 п.1/п.2): + + 1. `db.execute()` — блокирующий вызов. На loop'е он держит ВЕСЬ инстанс + (один воркер uvicorn, #3083) на время чекаута коннекта: при занятом пуле + это `pool_timeout` секунд, и никакой `asyncio.wait_for` его не прервёт. + В потоке ждёт поток, а loop продолжает обслуживать остальных. + 2. `commit()`/`rollback()` в конце — не про данные, а про КОННЕКТ: сессия + отдаёт его в пул только на завершении транзакции. Без этого коннект, + взятый ради одного SELECT'а кэша, живёт до конца функции — включая + ожидание чужого HTTP (замер 11.09: 8.5 с на догрузку Яндекса), и пиковый + спрос считается не в коротких SELECT'ах, а в целых фетчах. + + ЧТО ЭТО МЕНЯЕТ ДЛЯ ВЫЗЫВАЮЩЕГО: незакоммиченная работа его транзакции + (например fill-null backfill ФИАСа в `estimate_quality`) фиксируется здесь, + а не в первом commit'е ниже по коду. Новым классом поведения это не + является: обе функции-источника и так коммитят ЧУЖУЮ сессию сразу после + фетча (`save_imv_evaluation`, UPSERT в `external_valuations`) — сдвигается + момент, а не факт. Ошибка внутри `work` → `rollback` + проброс наверх, как + и раньше: гасят её существующие `except` вокруг вызова. + """ + + def _run() -> T: + try: + result = work() + except Exception: + db.rollback() + raise + db.commit() + return result + + return await asyncio.to_thread(_run) + + async def _get_or_fetch_imv_cached( db: Session, *, @@ -731,10 +768,12 @@ async def _get_or_fetch_imv_cached( has_loggia, ) - existing = ( - db.execute( - text( - """ + existing = await _db_step( + db, + lambda: ( + db.execute( + text( + """ SELECT id, cache_key, address, rooms, area_m2, floor, floor_at_home, house_type, renovation_type, has_balcony, has_loggia, lat, lon, geo_hash, avito_address_id, avito_location_id, @@ -747,11 +786,12 @@ async def _get_or_fetch_imv_cached( ORDER BY fetched_at DESC LIMIT 1 """ - ), - {"ck": cache_key, "ttl_hours": IMV_CACHE_TTL_HOURS}, - ) - .mappings() - .first() + ), + {"ck": cache_key, "ttl_hours": IMV_CACHE_TTL_HOURS}, + ) + .mappings() + .first() + ), ) if existing is not None: @@ -808,7 +848,9 @@ async def _get_or_fetch_imv_cached( # release в finally на всех выходах (исключение/таймаут — тоже). proxy_provider=RealProxyProvider(), ) - save_imv_evaluation(db, result, estimate_id=estimate_id_for_link) + await _db_step( + db, lambda: save_imv_evaluation(db, result, estimate_id=estimate_id_for_link) + ) logger.info( "imv: fresh recommended=%d range=(%d, %d) count=%d", result.recommended_price, @@ -842,7 +884,9 @@ async def _get_or_fetch_imv_cached( config=RealScraperConfig(), proxy_provider=RealProxyProvider(), # #3386, см. первый вызов выше ) - save_imv_evaluation(db, result, estimate_id=estimate_id_for_link) + await _db_step( + db, lambda: save_imv_evaluation(db, result, estimate_id=estimate_id_for_link) + ) logger.info( "imv: retry OK recommended=%d range=(%d, %d) count=%d", result.recommended_price, @@ -895,6 +939,10 @@ _DEFERRED_REFRESH_TASKS: set[asyncio.Task[None]] = set() # Подобран под текущий прод: 1 воркер uvicorn, mem_limit 768m, max_connections # 100 у Postgres. Догрузка почти всё время ждёт чужой HTTP, поэтому небольшой # потолок не занижает пропускную способность прогрева. +# +# Коннект к БД задача держит только на самих SELECT/UPSERT кэша, а не весь фетч +# (`_db_step`, #3408): до правки одна догрузка Яндекса удерживала соединение +# 8.5 с (замер на проде 11.09), то есть восемь таких задач выедали пул целиком. _MAX_DEFERRED_REFRESH_TASKS = 8 @@ -976,10 +1024,12 @@ async def _get_or_fetch_yandex_valuation_cached( # Cache lookup try: - cached = ( - db.execute( - text( - """ + cached = await _db_step( + db, + lambda: ( + db.execute( + text( + """ SELECT raw_payload, fetched_at FROM external_valuations WHERE source = 'yandex_valuation' @@ -988,11 +1038,12 @@ async def _get_or_fetch_yandex_valuation_cached( ORDER BY fetched_at DESC LIMIT 1 """ - ), - {"ck": cache_key}, - ) - .mappings() - .first() + ), + {"ck": cache_key}, + ) + .mappings() + .first() + ), ) except Exception as e: logger.warning("yandex_valuation: cache lookup failed: %s", e) @@ -1066,9 +1117,11 @@ async def _get_or_fetch_yandex_valuation_cached( # Save to cache (UPSERT on (source, cache_key)) try: - db.execute( - text( - """ + await _db_step( + db, + lambda: db.execute( + text( + """ INSERT INTO external_valuations ( source, cache_key, address, house_id, @@ -1088,16 +1141,16 @@ async def _get_or_fetch_yandex_valuation_cached( fetched_at = NOW(), expires_at = NOW() + (:ttl_hours || ' hours')::interval """ + ), + { + "ck": cache_key, + "addr": address, + "hid": house_id, + "payload": json.dumps(result.model_dump(mode="json"), ensure_ascii=False), + "ttl_hours": YANDEX_VALUATION_CACHE_TTL_HOURS, + }, ), - { - "ck": cache_key, - "addr": address, - "hid": house_id, - "payload": json.dumps(result.model_dump(mode="json"), ensure_ascii=False), - "ttl_hours": YANDEX_VALUATION_CACHE_TTL_HOURS, - }, ) - db.commit() logger.info( "yandex_valuation: fresh fetch saved key=%s items=%d", cache_key[:8], diff --git a/tradein-mvp/backend/tests/test_3408_estimator_db_off_loop.py b/tradein-mvp/backend/tests/test_3408_estimator_db_off_loop.py new file mode 100644 index 00000000..68db9a4a --- /dev/null +++ b/tradein-mvp/backend/tests/test_3408_estimator_db_off_loop.py @@ -0,0 +1,213 @@ +"""Синхронная БД внутри async-источников эстиматора: не на loop'е и не через фетч (#3408). + +`/estimate` публичен (meraocenka.ru) и крутится ОДНИМ воркером uvicorn (#3083), поэтому +у `db.execute()` прямо на event loop'е цена не «медленнее на миллисекунды», а «весь +инстанс не отвечает»: чекаут коннекта из занятого пула ждёт `pool_timeout`, и +`asyncio.wait_for` (`_with_budget`) синхронный вызов прервать не может. + +Что меряют тесты (значение, а не форму вызова): + (а) пока `db.execute` «ходит в БД» 0.3 с, соседняя корутина тикает — на origin/main + тиков ноль, loop стоит целиком (обе функции-источника: Яндекс и IMV); + (б) на время внешнего HTTP коннект ОТДАН пулу: пул из ОДНОГО коннекта, и второй + чекаут во время фетча обязан состояться. На origin/main транзакция, открытая + SELECT'ом кэша, живёт весь фетч (замер на проде 11.09: 8.5 с на одну догрузку + Яндекса при потолке 8 одновременных) — второй чекаут падает по таймауту. + +NB к (б): на sqlite SELECT кэша падает (Postgres-синтаксис `NOW()`), и это ровно тот +путь, что и на успехе: коннект удерживается транзакцией до `commit`/`rollback` +независимо от того, вернул ли SELECT строки. Проверяется «коннект свободен во время +await», а не «SELECT отработал». +""" + +from __future__ import annotations + +import asyncio +import os +import time +from typing import Any + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +import pytest +from sqlalchemy import create_engine, text +from sqlalchemy.exc import TimeoutError as SATimeoutError +from sqlalchemy.orm import sessionmaker +from sqlalchemy.pool import QueuePool + +from app.services import estimator as est + +# Столько «занимает поход в БД». 0.3 с — заметно больше планировочного шума и +# заметно меньше любого таймаута теста. +_DB_SLEEP_S = 0.3 + + +class _SlowSession: + """Session-дублёр: `execute` блокирует ПОТОК, как настоящий чекаут+SELECT.""" + + def __init__(self) -> None: + self.executed = 0 + self.commits = 0 + self.rollbacks = 0 + + def execute(self, *_a: Any, **_kw: Any) -> Any: + self.executed += 1 + time.sleep(_DB_SLEEP_S) + return self + + def mappings(self) -> Any: + return self + + def first(self) -> None: + return None + + def commit(self) -> None: + self.commits += 1 + + def rollback(self) -> None: + self.rollbacks += 1 + + +async def _ticks_during(call: Any) -> tuple[int, Any]: + """Сколько раз соседняя корутина успела проснуться, пока шёл `call`.""" + ticks = 0 + + async def _ticker() -> None: + nonlocal ticks + while True: + await asyncio.sleep(0) + ticks += 1 + + task = asyncio.create_task(_ticker()) + await asyncio.sleep(0) # дать тикеру стартовать + before = ticks + result = await call + during = ticks - before + task.cancel() + try: + await task + except asyncio.CancelledError: + pass + return during, result + + +# ── (а) sync-БД не держит event loop ──────────────────────────────────────── + + +async def test_yandex_cache_lookup_does_not_block_event_loop() -> None: + db = _SlowSession() + + during, result = await _ticks_during( + est._get_or_fetch_yandex_valuation_cached( + db, # type: ignore[arg-type] + address="Екатеринбург, улица Малышева, 84", + fetch_on_miss=False, # после промаха функция возвращает None сразу + ) + ) + + assert db.executed == 1, "SELECT кэша не звучал — тест ничего не проверил" + assert result is None + # Свободный loop успевает десятки тысяч тиков за 0.3 с; заблокированный — ноль. + assert during > 100, f"loop простоял весь SELECT кэша: тиков всего {during}" + + +async def test_imv_cache_lookup_does_not_block_event_loop(monkeypatch: Any) -> None: + db = _SlowSession() + + async def _no_fetch(**_kw: Any) -> None: + raise est.IMVTransientError("фетч в этом тесте не участвует") + + monkeypatch.setattr(est, "evaluate_via_imv", _no_fetch) + + during, result = await _ticks_during( + est._get_or_fetch_imv_cached( + db, # type: ignore[arg-type] + address="Екатеринбург, улица Малышева, 84", + rooms=2, + area_m2=54.0, + floor=3, + floor_at_home=9, + house_type="panel", + renovation_type="cosmetic", + has_balcony=True, + has_loggia=False, + ) + ) + + assert db.executed == 1, "SELECT кэша не звучал — тест ничего не проверил" + assert result is None + assert during > 100, f"loop простоял весь SELECT кэша: тиков всего {during}" + + +# ── (б) коннект не удерживается через внешний HTTP ────────────────────────── + + +class _FetchProbeScraper: + """Дублёр YandexValuationScraper: во время «фетча» пробует взять второй коннект.""" + + second_checkout_ok: bool | None = None + + def __init__(self, engine: Any) -> None: + self._engine = engine + + def __call__(self, *_a: Any, **_kw: Any) -> _FetchProbeScraper: + return self + + async def __aenter__(self) -> _FetchProbeScraper: + return self + + async def __aexit__(self, *_a: Any) -> None: + return None + + async def fetch_house_history(self, **_kw: Any) -> None: + await asyncio.sleep(0) + try: + with self._engine.connect() as conn: + conn.execute(text("SELECT 1")) + type(self).second_checkout_ok = True + except SATimeoutError: + type(self).second_checkout_ok = False + return None # «дом не найден» — функция выходит до UPSERT + + +@pytest.fixture +def one_connection_engine() -> Any: + """Пул ровно из одного коннекта: занятый коннект видно по таймауту чекаута.""" + engine = create_engine( + "sqlite://", + # sqlite по умолчанию берёт SingletonThreadPool — у него нет ни + # max_overflow, ни таймаута; нам нужен ровно тот пул, что на проде. + poolclass=QueuePool, + pool_size=1, + max_overflow=0, + pool_timeout=0.25, + connect_args={"check_same_thread": False}, + ) + try: + yield engine + finally: + engine.dispose() + + +async def test_yandex_fetch_does_not_hold_pool_connection( + monkeypatch: Any, one_connection_engine: Any +) -> None: + probe = _FetchProbeScraper(one_connection_engine) + _FetchProbeScraper.second_checkout_ok = None + monkeypatch.setattr(est, "YandexValuationScraper", probe) + + db = sessionmaker(bind=one_connection_engine, expire_on_commit=False)() + try: + result = await est._get_or_fetch_yandex_valuation_cached( + db, address="Екатеринбург, улица Малышева, 84" + ) + finally: + db.close() + + assert result is None + assert _FetchProbeScraper.second_checkout_ok is not None, ( + "фетч не звучал — тест ничего не проверил" + ) + assert _FetchProbeScraper.second_checkout_ok is True, ( + "коннект удерживается транзакцией всё время внешнего HTTP — " + "пик спроса на пул считается фетчами, а не SELECT'ами (#3408 п.1)" + ) diff --git a/tradein-mvp/backend/tests/test_3408_pool_ceiling.py b/tradein-mvp/backend/tests/test_3408_pool_ceiling.py new file mode 100644 index 00000000..9cbde0ef --- /dev/null +++ b/tradein-mvp/backend/tests/test_3408_pool_ceiling.py @@ -0,0 +1,54 @@ +"""Пул коннектов не меньше суммы потолков одновременности процесса (#3408 п.1). + +Процесс сам объявляет, сколько одновременной работы он допускает: 4 оценки, 4 +подсказки, 8 фоновых догрузок. Каждая из этих единиц работы держит СВОЮ сессию. +Если сумма больше пула, исчерпать пул можно штатной нагрузкой — и упереться не в +ресурс, а в `QueuePool limit ... timed out`, причём на публичной ручке. + +На origin/main сумма 16 против дефолта SQLAlchemy 5+10=15 — тест красный. + +Гейт нужен не ради текущих чисел, а ради будущих: поднять `_ESTIMATE_CONCURRENCY` +или `_MAX_DEFERRED_REFRESH_TASKS`, не поднимая пула, станет видно здесь. +Добавится воркер uvicorn (#3083) — числа per-process не меняются, но суммарный +потолок коннектов умножается на число воркеров; это отдельное решение, не это. +""" + +from __future__ import annotations + +import os + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +from app.api.public.mera import _SUGGEST_CONCURRENCY +from app.api.v1.trade_in import _ESTIMATE_CONCURRENCY +from app.core.db import engine +from app.services.estimator import _MAX_DEFERRED_REFRESH_TASKS + + +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 + + 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}с" + ) -- 2.45.3 From cf71825c273f67b8ab93ff7f55397743cc7a0cad Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 12 Sep 2026 02:20:26 +0500 Subject: [PATCH 2/3] =?UTF-8?q?fix(mera):=20=D0=BE=D1=82=D0=BC=D0=B5=D0=BD?= =?UTF-8?q?=D0=B0=20=D0=BF=D0=BE=20=D0=B1=D1=8E=D0=B4=D0=B6=D0=B5=D1=82?= =?UTF-8?q?=D1=83=20=D0=BE=D1=81=D1=82=D0=B0=D0=B2=D0=BB=D1=8F=D0=BB=D0=B0?= =?UTF-8?q?=20=D0=BE=D1=81=D0=B8=D1=80=D0=BE=D1=82=D0=B5=D0=B2=D1=88=D0=B8?= =?UTF-8?q?=D0=B9=20=D0=BF=D0=BE=D1=82=D0=BE=D0=BA=20=D0=B2=20=D1=87=D1=83?= =?UTF-8?q?=D0=B6=D0=BE=D0=B9=20Session?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Ревью PR #3444, M1. `_with_budget` — это `asyncio.wait_for`, а `asyncio.to_thread` отменить нельзя: снимается только ожидание со стороны loop'а. Корутина умирает, поток продолжает работать с ТОЙ ЖЕ `Session`, а вызывающий тем временем идёт дальше по своим шагам ПО ТОЙ ЖЕ сессии — следующий источник, `_fetch_anchor_comps`, `_persist_estimate_and_commit`. Два потока в одной сессии дают «another operation is in progress» / InvalidRequestError на следующем шаге БД: у источников её глушит `except` вокруг вызова, у персиста оценки не глушит ничего — 500 и потерянная оценка клиента, ровно под нагрузкой, ради которой PR и делается. `_db_step` теперь пробрасывает отмену ПОСЛЕ того, как поток отпустил сессию (`asyncio.shield` + ожидание шага). Цена — бюджет источника переезжает на длину ОДНОГО шага БД, а не на длину фетча, ради которой бюджет заведён. Почему не `threading.Lock` на сессию (вариант из ревью): лок внутри `_db_step` сериализует только шаги, которые через `_db_step` и проходят, — а названный пострадавший `_persist_estimate_and_commit` (estimator.py:5203) это ГОЛЫЙ `asyncio.to_thread(db...)`, как и ещё 17 мест эстиматора; лока они не берут, и дыра осталась бы открытой ровно там, где она стоит 500. Ожидание же в точке отмены закрывает ВСЕХ последующих потребителей сессии разом и не заводит глобального состояния (`WeakKeyDictionary`). Гейт по значению — tests/test_3408_db_step_cancel_orphan.py: следующий шаг (голый `to_thread`, как персист) не входит в сессию, пока сирота не закончил. Семантика проверена на питоне прода (3.12): `wait_for` по-прежнему отдаёт TimeoutError, источник деградирует в None. Остальное из ревью: - M2: комментарий у `_MAX_DEFERRED_REFRESH_TASKS` обещал за ОБА фоновых источника, а верен только для Яндекса. Циан держит коннект весь фетч (до 25 с): транзакцию открывают `_load_from_cache`/`load_session`, закрывает `db.commit()` в конце (scraper_kit .../cian/valuation.py:163,176,595). Формулировка сужена, остаток назван явно: функция общая со скраппером (cian_history_backfill.py:458), где коммит в середине менял бы семантику батча, — нужен отдельный опт-ин путь. На ПОТОЛОК пула остаток не влияет (коннект на задачу один независимо от того, как долго держится), только на среднюю занятость. - L1: `db.rollback()` после упавшего `_db_step` (estimator.py:1186) удалён — откат уже сделан в потоке, а на loop'е это блокирующий вызов. - L4: в core/db.py записано, что «пул >= суммы объявленных потолков» — ПОЛ, а не гарантия: коннект держит и любая ручка с `Depends(get_db)`, а глобального обработчика `sqlalchemy.exc.TimeoutError` в app/main.py нет (проверено: единственный handler — RequestValidationError, core/http_errors.py:59). - L3: гейт пула больше не читает `pool._timeout` и не молчит при переименовании `_max_overflow` — публичный `pool.size()` + приватное поле за явным assert'ом. `pool_timeout` из этого коммита УБРАН намеренно: это единственная правка, которая меняет режим отказа с «медленно» на «быстро с ошибкой», и она едет во все сервисы образа (backend, scraper, tgbot). Возвращается отдельным коммитом в конце ветки — чтобы ветку можно было смержить без него или откатить одной командой. Refs #3083, #3408 --- tradein-mvp/backend/app/core/db.py | 14 ++-- tradein-mvp/backend/app/services/estimator.py | 49 +++++++++-- .../tests/test_3408_db_step_cancel_orphan.py | 82 +++++++++++++++++++ .../backend/tests/test_3408_pool_ceiling.py | 35 ++++---- 4 files changed, 147 insertions(+), 33 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_3408_db_step_cancel_orphan.py 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}с" - ) -- 2.45.3 From ae6d28d5e2a315fa55898683340c5ebe01187c03 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 12 Sep 2026 02:23:09 +0500 Subject: [PATCH 3/3] =?UTF-8?q?fix(mera):=20pool=5Ftimeout=2030=E2=86=925?= =?UTF-8?q?=20=D1=81=20=E2=80=94=20=D0=BE=D1=82=D0=B4=D0=B5=D0=BB=D1=8C?= =?UTF-8?q?=D0=BD=D1=8B=D0=BC=20=D0=BA=D0=BE=D0=BC=D0=BC=D0=B8=D1=82=D0=BE?= =?UTF-8?q?=D0=BC,=20=D1=81=20=D1=82=D1=80=D0=B8=D0=B3=D0=B3=D0=B5=D1=80?= =?UTF-8?q?=D0=BE=D0=BC=20=D0=BE=D1=82=D0=BA=D0=B0=D1=82=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Единственная правка ветки, которая меняет РЕЖИМ ОТКАЗА при исчерпании пула: было «медленно» (ждём коннект до 30 с), стало «быстро с ошибкой» (5 с и `sqlalchemy.exc.TimeoutError` → 500, глобального обработчика в app/main.py нет). И едет она во все сервисы образа — backend, scraper, tgbot (tradein-mvp/docker-compose.prod.yml), для скраппера и бота обоснования в коде нет: за 29 ч логов исчерпания пула не было ни разу, проверить новое значение на проде пока не на чем. Поэтому коммит последний в ветке: ветку можно мержить без него, а на проде — откатить одной командой (`git revert`). Обоснование самого значения: чекаут коннекта нельзя прервать `asyncio.wait_for`, он занимает поток `asyncio.to_thread` целиком, а пул потоков конечен (min(32, cpu+4)) — исчерпанный пул коннектов превращается в исчерпанный пул потоков. 5 с короче самого короткого бюджета источника (8 с Yandex/Cian/ house_meta; geocode 12 с, IMV 20 с — длиннее): занятый пул деградирует ОДИН источник, а не весь запрос. ТРИГГЕР ОТКАТА (вернуть 30 с) записан в комментарии рядом со значением: любое `QueuePool limit ... timed out` в логах бэкенда ЛИБО рост failed+zombie в `scrape_runs` после деплоя. Правка комментария по ревью (L2): «вчетверо больше любого бюджета внешнего источника (8 с)» было неточно — бюджеты 8 / 12 / 20 с, перечислены явно. Гейт `test_pool_checkout_wait_shorter_than_source_budget` переехал сюда же (в коммите без `pool_timeout` он был бы красным) и читает публичный `engine.pool.timeout()` вместо приватного `pool._timeout`. Refs #3083, #3408 --- tradein-mvp/backend/app/core/db.py | 16 +++++++++++++++ .../backend/tests/test_3408_pool_ceiling.py | 20 +++++++++++++++++++ 2 files changed, 36 insertions(+) diff --git a/tradein-mvp/backend/app/core/db.py b/tradein-mvp/backend/app/core/db.py index 408d47ba..1214a3eb 100644 --- a/tradein-mvp/backend/app/core/db.py +++ b/tradein-mvp/backend/app/core/db.py @@ -35,6 +35,22 @@ engine = create_engine( # которые пул закрывает сам. Потолок процесса: 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) diff --git a/tradein-mvp/backend/tests/test_3408_pool_ceiling.py b/tradein-mvp/backend/tests/test_3408_pool_ceiling.py index 254340db..2c220834 100644 --- a/tradein-mvp/backend/tests/test_3408_pool_ceiling.py +++ b/tradein-mvp/backend/tests/test_3408_pool_ceiling.py @@ -49,3 +49,23 @@ def test_pool_ceiling_covers_declared_concurrency() -> None: f"({_ESTIMATE_CONCURRENCY} оценок + {_SUGGEST_CONCURRENCY} подсказок + " f"{_MAX_DEFERRED_REFRESH_TASKS} фоновых догрузок) — штатная работа исчерпает пул" ) + + +def test_pool_checkout_wait_shorter_than_source_budget() -> None: + """Ожидание коннекта короче бюджета внешнего источника. + + Иначе занятый пул съедает весь бюджет запроса (и поток `asyncio.to_thread`, + которых тоже конечное число) вместо того, чтобы деградировать один источник. + Сравниваем с самым КОРОТКИМ бюджетом (8 с Yandex/Cian): geocode 12 с и IMV + 20 с длиннее, их этот же потолок покрывает с запасом. + """ + 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}с" + ) -- 2.45.3