Compare commits

..

4 commits

Author SHA1 Message Date
4efcb712e4 Merge pull request 'fix(mera/estimate): коннект БД не живёт через внешний HTTP; потолок пула ≥ суммы потолков одновременности (#3083, #3408)' (#3444) from fix/3083-estimate-throughput into main
All checks were successful
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 4m14s
Deploy Trade-In / build-backend (push) Successful in 1m10s
Deploy Trade-In / deploy (push) Successful in 2m44s
Deploy Trade-In / deploy-status (push) Successful in 2s
Deploy Trade-In / perimeter-smoke (push) Successful in 1m41s
2026-09-11 21:33:14 +00:00
ae6d28d5e2 fix(mera): pool_timeout 30→5 с — отдельным коммитом, с триггером отката
All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m4s
Единственная правка ветки, которая меняет РЕЖИМ ОТКАЗА при исчерпании пула:
было «медленно» (ждём коннект до 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
2026-09-12 02:23:09 +05:00
cf71825c27 fix(mera): отмена по бюджету оставляла осиротевший поток в чужой Session
Ревью 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
2026-09-12 02:20:26 +05:00
9f696299de fix(mera): sync-БД источников эстиматора — с event loop в поток и не через фетч
All checks were successful
CI / changes (pull_request) Successful in 10s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI Trade-In / changes (pull_request) Successful in 9s
CI / backend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m16s
Замер на проде 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
2026-09-12 01:01:54 +05:00
6 changed files with 525 additions and 34 deletions

View file

@ -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).

View file

@ -16,6 +16,41 @@ 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.
#
# Инвариант «пул >= суммы объявленных потолков» — это ПОЛ, а не гарантия:
# сумма считает по одному коннекту на объявленную единицу работы, а коннект
# держит и всё остальное — любая ручка с `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)

View file

@ -695,6 +695,68 @@ 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'ах, а в целых фетчах.
3. Отмена шага ДОЖИДАЕТСЯ потока (`asyncio.shield` ниже). Поток отменить
нельзя, а сессия у шага не своя общая со всем остальным запросом.
ЧТО ЭТО МЕНЯЕТ ДЛЯ ВЫЗЫВАЮЩЕГО: незакоммиченная работа его транзакции
(например 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
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(
db: Session,
*,
@ -731,10 +793,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 +811,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 +873,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 +909,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 +964,19 @@ _DEFERRED_REFRESH_TASKS: set[asyncio.Task[None]] = set()
# Подобран под текущий прод: 1 воркер uvicorn, mem_limit 768m, max_connections
# 100 у Postgres. Догрузка почти всё время ждёт чужой HTTP, поэтому небольшой
# потолок не занижает пропускную способность прогрева.
#
# Фоновых источника ДВА, и по удержанию коннекта они РАЗНЫЕ (#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
@ -976,10 +1058,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 +1072,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 +1151,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,24 +1175,25 @@ 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],
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

View file

@ -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, "сирота не завершил транзакцию — коннект не вернулся в пул"

View file

@ -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)"
)

View file

@ -0,0 +1,71 @@
"""Пул коннектов не меньше суммы потолков одновременности процесса (#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
from typing import Any
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 _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() + _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`,
которых тоже конечное число) вместо того, чтобы деградировать один источник.
Сравниваем с самым КОРОТКИМ бюджетом (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}с"
)