All checks were successful
CI Trade-In / 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 / changes (pull_request) Successful in 10s
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 5m24s
`asyncio.to_thread` отменить нельзя: по истечении бюджета (`_with_budget` = `asyncio.wait_for`, у геокодера 12 с) снимается только ожидание со стороны loop'а — поток продолжает работать с ТОЙ ЖЕ `Session`, что и весь запрос. Вызывающий тем временем идёт дальше: следующий источник, `_fetch_anchor_comps`, `_persist_estimate_and_commit`. Два потока в одной `Session` дают «another operation is in progress» / InvalidRequestError на СЛЕДУЮЩЕМ шаге. У источников эту ошибку глушит `except` вокруг вызова, у персиста оценки не глушит никто — 500 и потерянная оценка клиента. `app/core/db.py: run_db_thread` — ТОЛЬКО защита от сироты: `ensure_future` + `shield`, на отмене дождаться потока (`asyncio.wait`), прочитать `step.exception()` (иначе asyncio печатает «Task exception was never retrieved» без контекста) и пробросить отмену. Commit/rollback туда НЕ вынесены: посреди геокодинга commit зафиксировал бы частичное состояние оценки. `estimator._db_step` переписан поверх и добавляет свои commit/rollback сам — его поведение не меняется, гейт tests/test_3408_db_step_cancel_orphan.py остаётся зелёным. Заменено 34 вызова, работающих по сессии запроса: 12 в geocoder.py (кэш-чтение и записи, геопортал, кадастр, houses, reverse, suggest), 19 в estimator.py (в т.ч. `_backfill_house_fias`, `_save_yandex_history_items`, `_fetch_anchor_comps`, `_price_from_inputs` с db-резолверами, персист оценки, `_fetch_price_trend`, `_is_premium_building`), 2 в api/v1/geocode.py, 1 в api/v1/privacy_admin.py. Не тронуты вызовы со СВОЕЙ сессией: `user_events.schedule_event` (внутри `record_event` свой `SessionLocal`) и `sber_index` (сессия задачи планировщика, отменять её некому). Гейт по значению — tests/test_3449_geocoder_cancel_orphan.py: отмена по бюджету во время шага БД геокодера, следом ГОЛЫЙ `to_thread(db.execute, ...)` (образец персиста); проверяется, что он не вошёл в сессию, пока сирота ещё в ней. На исходном коде тест краснеет: conflicts == 1. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
108 lines
7.5 KiB
Python
108 lines
7.5 KiB
Python
import asyncio
|
||
import logging
|
||
from collections.abc import Callable, Generator
|
||
from typing import Any
|
||
|
||
from sqlalchemy import create_engine
|
||
from sqlalchemy.orm import DeclarativeBase, Session, sessionmaker
|
||
|
||
from app.core.config import settings
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
engine = create_engine(
|
||
settings.database_url,
|
||
pool_pre_ping=True,
|
||
future=True,
|
||
# #3194: SQLAlchemy печатает ВСЕ bind-параметры в тексте StatementError —
|
||
# через них в GlitchTip уезжали ключ шифрования кук и сами куки
|
||
# (pgp_sym_encrypt(:cookies_json, :key)). Флаг на УРОВНЕ ДВИЖКА кроет все
|
||
# сайты вызова разом, включая будущие.
|
||
# НЕ закрывает: текст ошибки самого драйвера (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)
|
||
|
||
|
||
class Base(DeclarativeBase):
|
||
pass
|
||
|
||
|
||
def get_db() -> Generator[Session, None, None]:
|
||
db = SessionLocal()
|
||
try:
|
||
yield db
|
||
finally:
|
||
db.close()
|
||
|
||
|
||
async def run_db_thread[T](fn: Callable[..., T], *args: Any, **kwargs: Any) -> T:
|
||
"""`asyncio.to_thread(fn, ...)`, который при отмене ДОЖИДАЕТСЯ своего потока (#3449).
|
||
|
||
Для синхронной работы по сессии, которая шагу НЕ принадлежит — она общая со
|
||
всем остальным запросом (`Depends(get_db)`). Поток отменить нельзя: у
|
||
`asyncio.to_thread` отменяется только ожидание со стороны loop'а. Корутина
|
||
умирает по бюджету источника (`estimator._with_budget` = `asyncio.wait_for`,
|
||
геокодер — 12 с), а поток продолжает работать с ТОЙ ЖЕ `Session`, пока
|
||
вызывающий уже идёт дальше по коду — следующий источник,
|
||
`_fetch_anchor_comps`, `_persist_estimate_and_commit`. Два потока в одной
|
||
`Session` дают «another operation is in progress» / `InvalidRequestError` на
|
||
СЛЕДУЮЩЕМ шаге: у источников такую ошибку глушит `except` вокруг вызова, у
|
||
персиста оценки не глушит ничего — 500 и потерянная оценка клиента.
|
||
|
||
Поэтому отмена пробрасывается ПОСЛЕ того, как поток отпустил сессию. Цена —
|
||
бюджет источника переезжает на длину ОДНОГО шага БД (чекаут ≤ `pool_timeout`
|
||
плюс сам запрос), а не на длину фетча, ради которой бюджет и заведён.
|
||
|
||
Только защита от сироты: транзакцию шаг НЕ завершает (в середине геокодинга
|
||
commit зафиксировал бы частичное состояние оценки). Кому нужен ещё и возврат
|
||
коннекта в пул перед внешним HTTP — `estimator._db_step`, он поверх этого.
|
||
|
||
Гейт — tests/test_3449_geocoder_cancel_orphan.py.
|
||
"""
|
||
step = asyncio.ensure_future(asyncio.to_thread(fn, *args, **kwargs))
|
||
try:
|
||
return await asyncio.shield(step)
|
||
except asyncio.CancelledError:
|
||
await asyncio.wait([step])
|
||
if not step.cancelled() and step.exception() is not None:
|
||
# Результата уже никто не ждёт: без явного чтения asyncio напечатает
|
||
# «Task exception was never retrieved» вообще без контекста.
|
||
logger.warning("шаг БД упал уже после отмены: %s", step.exception())
|
||
raise
|