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
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
This commit is contained in:
parent
a659b18771
commit
9f696299de
5 changed files with 374 additions and 33 deletions
|
|
@ -79,8 +79,10 @@ _estimate_limiter = SlidingWindowLimiter(
|
||||||
# внешние тиры; пила одновременных оценок выедает пул и тормозит весь /api/v1/*.
|
# внешние тиры; пила одновременных оценок выедает пул и тормозит весь /api/v1/*.
|
||||||
# Образец — public/mera.py::_suggest_slots (4 слота на секундное автодополнение).
|
# Образец — public/mera.py::_suggest_slots (4 слота на секундное автодополнение).
|
||||||
#
|
#
|
||||||
# 4 слота: вместе с 4 слотами suggest — 8 одновременно удерживаемых соединений
|
# 4 слота: вместе с 4 слотами suggest — 8 одновременно удерживаемых соединений.
|
||||||
# из 15 возможных, остаток пула остаётся прочим ручкам. Ожидание слота 5с ≈ две
|
# Пул под это заведомо шире: 5+15=20 на процесс, и потолок пула держится не
|
||||||
|
# меньше СУММЫ объявленных потолков одновременности, включая 8 фоновых догрузок
|
||||||
|
# (core/db.py + tests/test_3408_pool_ceiling.py, #3408). Ожидание слота 5с ≈ две
|
||||||
# длительности оценки: если за это время слот не освободился, очередь глубока и
|
# длительности оценки: если за это время слот не освободился, очередь глубока и
|
||||||
# честный ответ — быстрый 429 с Retry-After, а не растущая очередь (очередь под
|
# честный ответ — быстрый 429 с Retry-After, а не растущая очередь (очередь под
|
||||||
# нагрузкой — те же занятые соединения плюс таймаут у клиента; mera.py:117-127).
|
# нагрузкой — те же занятые соединения плюс таймаут у клиента; mera.py:117-127).
|
||||||
|
|
|
||||||
|
|
@ -16,6 +16,25 @@ engine = create_engine(
|
||||||
# НЕ закрывает: текст ошибки самого драйвера (Postgres DETAIL со значением)
|
# НЕ закрывает: текст ошибки самого драйвера (Postgres DETAIL со значением)
|
||||||
# и сырые psycopg-подключения мимо движков — это отдельный класс.
|
# и сырые psycopg-подключения мимо движков — это отдельный класс.
|
||||||
hide_parameters=True,
|
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)
|
SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine, expire_on_commit=False)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -695,6 +695,43 @@ YANDEX_VALUATION_DEFAULT_CATEGORY = "APARTMENT"
|
||||||
YANDEX_VALUATION_DEFAULT_TYPE = "SELL"
|
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(
|
async def _get_or_fetch_imv_cached(
|
||||||
db: Session,
|
db: Session,
|
||||||
*,
|
*,
|
||||||
|
|
@ -731,10 +768,12 @@ async def _get_or_fetch_imv_cached(
|
||||||
has_loggia,
|
has_loggia,
|
||||||
)
|
)
|
||||||
|
|
||||||
existing = (
|
existing = await _db_step(
|
||||||
db.execute(
|
db,
|
||||||
text(
|
lambda: (
|
||||||
"""
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
SELECT id, cache_key, address, rooms, area_m2, floor, floor_at_home,
|
SELECT id, cache_key, address, rooms, area_m2, floor, floor_at_home,
|
||||||
house_type, renovation_type, has_balcony, has_loggia,
|
house_type, renovation_type, has_balcony, has_loggia,
|
||||||
lat, lon, geo_hash, avito_address_id, avito_location_id,
|
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
|
ORDER BY fetched_at DESC
|
||||||
LIMIT 1
|
LIMIT 1
|
||||||
"""
|
"""
|
||||||
),
|
),
|
||||||
{"ck": cache_key, "ttl_hours": IMV_CACHE_TTL_HOURS},
|
{"ck": cache_key, "ttl_hours": IMV_CACHE_TTL_HOURS},
|
||||||
)
|
)
|
||||||
.mappings()
|
.mappings()
|
||||||
.first()
|
.first()
|
||||||
|
),
|
||||||
)
|
)
|
||||||
|
|
||||||
if existing is not None:
|
if existing is not None:
|
||||||
|
|
@ -808,7 +848,9 @@ async def _get_or_fetch_imv_cached(
|
||||||
# release в finally на всех выходах (исключение/таймаут — тоже).
|
# release в finally на всех выходах (исключение/таймаут — тоже).
|
||||||
proxy_provider=RealProxyProvider(),
|
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(
|
logger.info(
|
||||||
"imv: fresh recommended=%d range=(%d, %d) count=%d",
|
"imv: fresh recommended=%d range=(%d, %d) count=%d",
|
||||||
result.recommended_price,
|
result.recommended_price,
|
||||||
|
|
@ -842,7 +884,9 @@ async def _get_or_fetch_imv_cached(
|
||||||
config=RealScraperConfig(),
|
config=RealScraperConfig(),
|
||||||
proxy_provider=RealProxyProvider(), # #3386, см. первый вызов выше
|
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(
|
logger.info(
|
||||||
"imv: retry OK recommended=%d range=(%d, %d) count=%d",
|
"imv: retry OK recommended=%d range=(%d, %d) count=%d",
|
||||||
result.recommended_price,
|
result.recommended_price,
|
||||||
|
|
@ -895,6 +939,10 @@ _DEFERRED_REFRESH_TASKS: set[asyncio.Task[None]] = set()
|
||||||
# Подобран под текущий прод: 1 воркер uvicorn, mem_limit 768m, max_connections
|
# Подобран под текущий прод: 1 воркер uvicorn, mem_limit 768m, max_connections
|
||||||
# 100 у Postgres. Догрузка почти всё время ждёт чужой HTTP, поэтому небольшой
|
# 100 у Postgres. Догрузка почти всё время ждёт чужой HTTP, поэтому небольшой
|
||||||
# потолок не занижает пропускную способность прогрева.
|
# потолок не занижает пропускную способность прогрева.
|
||||||
|
#
|
||||||
|
# Коннект к БД задача держит только на самих SELECT/UPSERT кэша, а не весь фетч
|
||||||
|
# (`_db_step`, #3408): до правки одна догрузка Яндекса удерживала соединение
|
||||||
|
# 8.5 с (замер на проде 11.09), то есть восемь таких задач выедали пул целиком.
|
||||||
_MAX_DEFERRED_REFRESH_TASKS = 8
|
_MAX_DEFERRED_REFRESH_TASKS = 8
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -976,10 +1024,12 @@ async def _get_or_fetch_yandex_valuation_cached(
|
||||||
|
|
||||||
# Cache lookup
|
# Cache lookup
|
||||||
try:
|
try:
|
||||||
cached = (
|
cached = await _db_step(
|
||||||
db.execute(
|
db,
|
||||||
text(
|
lambda: (
|
||||||
"""
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
SELECT raw_payload, fetched_at
|
SELECT raw_payload, fetched_at
|
||||||
FROM external_valuations
|
FROM external_valuations
|
||||||
WHERE source = 'yandex_valuation'
|
WHERE source = 'yandex_valuation'
|
||||||
|
|
@ -988,11 +1038,12 @@ async def _get_or_fetch_yandex_valuation_cached(
|
||||||
ORDER BY fetched_at DESC
|
ORDER BY fetched_at DESC
|
||||||
LIMIT 1
|
LIMIT 1
|
||||||
"""
|
"""
|
||||||
),
|
),
|
||||||
{"ck": cache_key},
|
{"ck": cache_key},
|
||||||
)
|
)
|
||||||
.mappings()
|
.mappings()
|
||||||
.first()
|
.first()
|
||||||
|
),
|
||||||
)
|
)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning("yandex_valuation: cache lookup failed: %s", 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))
|
# Save to cache (UPSERT on (source, cache_key))
|
||||||
try:
|
try:
|
||||||
db.execute(
|
await _db_step(
|
||||||
text(
|
db,
|
||||||
"""
|
lambda: db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
INSERT INTO external_valuations (
|
INSERT INTO external_valuations (
|
||||||
source, cache_key, address,
|
source, cache_key, address,
|
||||||
house_id,
|
house_id,
|
||||||
|
|
@ -1088,16 +1141,16 @@ async def _get_or_fetch_yandex_valuation_cached(
|
||||||
fetched_at = NOW(),
|
fetched_at = NOW(),
|
||||||
expires_at = NOW() + (:ttl_hours || ' hours')::interval
|
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(
|
logger.info(
|
||||||
"yandex_valuation: fresh fetch saved key=%s items=%d",
|
"yandex_valuation: fresh fetch saved key=%s items=%d",
|
||||||
cache_key[:8],
|
cache_key[:8],
|
||||||
|
|
|
||||||
213
tradein-mvp/backend/tests/test_3408_estimator_db_off_loop.py
Normal file
213
tradein-mvp/backend/tests/test_3408_estimator_db_off_loop.py
Normal 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)"
|
||||||
|
)
|
||||||
54
tradein-mvp/backend/tests/test_3408_pool_ceiling.py
Normal file
54
tradein-mvp/backend/tests/test_3408_pool_ceiling.py
Normal file
|
|
@ -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}с"
|
||||||
|
)
|
||||||
Loading…
Add table
Reference in a new issue