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
This commit is contained in:
bot-backend 2026-09-12 02:20:26 +05:00
parent 9f696299de
commit cf71825c27
4 changed files with 147 additions and 33 deletions

View file

@ -23,18 +23,18 @@ engine = create_engine(
# Дефолт SQLAlchemy 5+10=15 меньше этой суммы, то есть исчерпать пул можно # Дефолт SQLAlchemy 5+10=15 меньше этой суммы, то есть исчерпать пул можно
# штатной работой, не абузом. Гейт — tests/test_3408_pool_ceiling.py. # штатной работой, не абузом. Гейт — tests/test_3408_pool_ceiling.py.
# #
# Инвариант «пул >= суммы объявленных потолков» — это ПОЛ, а не гарантия:
# сумма считает по одному коннекту на объявленную единицу работы, а коннект
# держит и всё остальное — любая ручка с `Depends(get_db)` тоже. Глобального
# обработчика `sqlalchemy.exc.TimeoutError` в app/main.py нет, поэтому на
# исчерпанном пуле соседние ручки отдают 500 (после ожидания `pool_timeout`),
# а не медленный 200.
#
# pool_size оставлен дефолтным (5): это ПОСТОЯННО открытые коннекты, а в покое # pool_size оставлен дефолтным (5): это ПОСТОЯННО открытые коннекты, а в покое
# прод держит 5-6 (замер 11.09). Растёт только overflow — коннекты пика, # прод держит 5-6 (замер 11.09). Растёт только overflow — коннекты пика,
# которые пул закрывает сам. Потолок процесса: 5 + 15 = 20; воркер один # которые пул закрывает сам. Потолок процесса: 5 + 15 = 20; воркер один
# (docker-compose.prod.yml, uvicorn без --workers), Postgres max_connections=100. # (docker-compose.prod.yml, uvicorn без --workers), Postgres max_connections=100.
max_overflow=15, 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)

View file

@ -696,9 +696,9 @@ YANDEX_VALUATION_DEFAULT_TYPE = "SELL"
async def _db_step[T](db: Session, work: Callable[[], T]) -> T: async def _db_step[T](db: Session, work: Callable[[], T]) -> T:
"""Синхронный шаг БД внутри async-функции: в потоке И с возвратом коннекта в пул. """Синхронный шаг БД внутри async-функции: в потоке, с возвратом коннекта и без сирот.
Две вещи сразу, потому что обе про одно и то же `Session` из `get_db` Три вещи сразу, потому что все про одно и то же `Session` из `get_db`
(#3408 п.1/п.2): (#3408 п.1/п.2):
1. `db.execute()` блокирующий вызов. На loop'е он держит ВЕСЬ инстанс 1. `db.execute()` блокирующий вызов. На loop'е он держит ВЕСЬ инстанс
@ -710,6 +710,8 @@ async def _db_step[T](db: Session, work: Callable[[], T]) -> T:
взятый ради одного SELECT'а кэша, живёт до конца функции — включая взятый ради одного SELECT'а кэша, живёт до конца функции — включая
ожидание чужого HTTP (замер 11.09: 8.5 с на догрузку Яндекса), и пиковый ожидание чужого HTTP (замер 11.09: 8.5 с на догрузку Яндекса), и пиковый
спрос считается не в коротких SELECT'ах, а в целых фетчах. спрос считается не в коротких SELECT'ах, а в целых фетчах.
3. Отмена шага ДОЖИДАЕТСЯ потока (`asyncio.shield` ниже). Поток отменить
нельзя, а сессия у шага не своя общая со всем остальным запросом.
ЧТО ЭТО МЕНЯЕТ ДЛЯ ВЫЗЫВАЮЩЕГО: незакоммиченная работа его транзакции ЧТО ЭТО МЕНЯЕТ ДЛЯ ВЫЗЫВАЮЩЕГО: незакоммиченная работа его транзакции
(например fill-null backfill ФИАСа в `estimate_quality`) фиксируется здесь, (например fill-null backfill ФИАСа в `estimate_quality`) фиксируется здесь,
@ -729,7 +731,30 @@ async def _db_step[T](db: Session, work: Callable[[], T]) -> T:
db.commit() db.commit()
return result 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( async def _get_or_fetch_imv_cached(
@ -940,9 +965,18 @@ _DEFERRED_REFRESH_TASKS: set[asyncio.Task[None]] = set()
# 100 у Postgres. Догрузка почти всё время ждёт чужой HTTP, поэтому небольшой # 100 у Postgres. Догрузка почти всё время ждёт чужой HTTP, поэтому небольшой
# потолок не занижает пропускную способность прогрева. # потолок не занижает пропускную способность прогрева.
# #
# Коннект к БД задача держит только на самих SELECT/UPSERT кэша, а не весь фетч # Фоновых источника ДВА, и по удержанию коннекта они РАЗНЫЕ (#3408):
# (`_db_step`, #3408): до правки одна догрузка Яндекса удерживала соединение # • yandex_valuation — коннект только на самих SELECT/UPSERT кэша (`_db_step`):
# 8.5 с (замер на проде 11.09), то есть восемь таких задач выедали пул целиком. # до правки одна догрузка держала соединение 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 _MAX_DEFERRED_REFRESH_TASKS = 8
@ -1157,8 +1191,9 @@ async def _get_or_fetch_yandex_valuation_cached(
len(result.history_items), len(result.history_items),
) )
except Exception as e: except Exception as e:
# Откат уже сделан в потоке (`_db_step` → `except` → `db.rollback()`), а здесь
# это был бы блокирующий вызов на event loop'е — ровно то, что правка убирает.
logger.warning("yandex_valuation: cache save failed (continuing): %s", e) logger.warning("yandex_valuation: cache save failed (continuing): %s", e)
db.rollback()
return result 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

@ -16,6 +16,7 @@
from __future__ import annotations from __future__ import annotations
import os import os
from typing import Any
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") 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 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: def test_pool_ceiling_covers_declared_concurrency() -> None:
declared = _ESTIMATE_CONCURRENCY + _SUGGEST_CONCURRENCY + _MAX_DEFERRED_REFRESH_TASKS 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, ( assert ceiling >= declared, (
f"пул {ceiling} меньше суммы потолков одновременности {declared} " f"пул {ceiling} меньше суммы потолков одновременности {declared} "
f"({_ESTIMATE_CONCURRENCY} оценок + {_SUGGEST_CONCURRENCY} подсказок + " f"({_ESTIMATE_CONCURRENCY} оценок + {_SUGGEST_CONCURRENCY} подсказок + "
f"{_MAX_DEFERRED_REFRESH_TASKS} фоновых догрузок) — штатная работа исчерпает пул" 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}с"
)