From c251c02f1ec846f0d749e4a0bfc7409d59113126 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 12 Sep 2026 13:20:44 +0500 Subject: [PATCH 1/2] =?UTF-8?q?=D0=9E=D1=82=D0=BC=D0=B5=D0=BD=D0=B0=20?= =?UTF-8?q?=D0=BF=D0=BE=20=D0=B1=D1=8E=D0=B4=D0=B6=D0=B5=D1=82=D1=83=20?= =?UTF-8?q?=D0=B1=D0=BE=D0=BB=D1=8C=D1=88=D0=B5=20=D0=BD=D0=B5=20=D0=BE?= =?UTF-8?q?=D1=81=D1=82=D0=B0=D0=B2=D0=BB=D1=8F=D0=B5=D1=82=20=D1=81=D0=B8?= =?UTF-8?q?=D1=80=D0=BE=D1=82=D1=83=20=D0=B2=20=D1=81=D0=B5=D1=81=D1=81?= =?UTF-8?q?=D0=B8=D0=B8=20=D0=B7=D0=B0=D0=BF=D1=80=D0=BE=D1=81=D0=B0=20(#3?= =?UTF-8?q?449)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `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 --- tradein-mvp/backend/app/api/v1/geocode.py | 7 +- .../backend/app/api/v1/privacy_admin.py | 5 +- tradein-mvp/backend/app/core/db.py | 43 ++++++++- tradein-mvp/backend/app/services/estimator.py | 70 +++++---------- tradein-mvp/backend/app/services/geocoder.py | 27 +++--- .../tests/test_3449_geocoder_cancel_orphan.py | 88 +++++++++++++++++++ 6 files changed, 172 insertions(+), 68 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_3449_geocoder_cancel_orphan.py diff --git a/tradein-mvp/backend/app/api/v1/geocode.py b/tradein-mvp/backend/app/api/v1/geocode.py index c1f18b9f..ae439589 100644 --- a/tradein-mvp/backend/app/api/v1/geocode.py +++ b/tradein-mvp/backend/app/api/v1/geocode.py @@ -2,7 +2,6 @@ from __future__ import annotations -import asyncio import logging from typing import Annotated @@ -11,7 +10,7 @@ from pydantic import BaseModel, Field from sqlalchemy import text from sqlalchemy.orm import Session -from app.core.db import get_db +from app.core.db import get_db, run_db_thread from app.services.estimator import _lookup_house_facts from app.services.geocoder import GeocodeResult, geocode, reverse_geocode, suggest @@ -254,9 +253,9 @@ async def house_facts( """ target_house_id: int | None = None if fias_id is not None: - target_house_id = await asyncio.to_thread(_resolve_house_id_by_fias, db, fias_id) + target_house_id = await run_db_thread(_resolve_house_id_by_fias, db, fias_id) - facts = await asyncio.to_thread( + facts = await run_db_thread( _lookup_house_facts, db, target_house_id=target_house_id, diff --git a/tradein-mvp/backend/app/api/v1/privacy_admin.py b/tradein-mvp/backend/app/api/v1/privacy_admin.py index 9ea0ee9d..85705a00 100644 --- a/tradein-mvp/backend/app/api/v1/privacy_admin.py +++ b/tradein-mvp/backend/app/api/v1/privacy_admin.py @@ -14,7 +14,6 @@ module docstring (честно про то, что не всегда разре from __future__ import annotations -import asyncio import logging from typing import Annotated from uuid import UUID @@ -23,7 +22,7 @@ from fastapi import APIRouter, Depends, HTTPException from pydantic import BaseModel, Field from sqlalchemy.orm import Session -from app.core.db import get_db +from app.core.db import get_db, run_db_thread logger = logging.getLogger(__name__) @@ -67,7 +66,7 @@ async def erase_person_data_endpoint( from app.services.data_erasure import erase_person_data - counters = await asyncio.to_thread( + counters = await run_db_thread( erase_person_data, db, username=payload.username, diff --git a/tradein-mvp/backend/app/core/db.py b/tradein-mvp/backend/app/core/db.py index 1214a3eb..94555e63 100644 --- a/tradein-mvp/backend/app/core/db.py +++ b/tradein-mvp/backend/app/core/db.py @@ -1,10 +1,15 @@ -from collections.abc import Generator +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, @@ -65,3 +70,39 @@ def get_db() -> Generator[Session, None, None]: 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 diff --git a/tradein-mvp/backend/app/services/estimator.py b/tradein-mvp/backend/app/services/estimator.py index 8556f6ec..87426284 100644 --- a/tradein-mvp/backend/app/services/estimator.py +++ b/tradein-mvp/backend/app/services/estimator.py @@ -56,7 +56,7 @@ from sqlalchemy import text from sqlalchemy.orm import Session from app.core.config import LISTINGS_FRESH_DAYS, Settings, settings -from app.core.db import SessionLocal +from app.core.db import SessionLocal, run_db_thread from app.schemas.trade_in import ( AggregatedEstimate, AnalogLot, @@ -803,8 +803,9 @@ async def _db_step[T](db: Session, work: Callable[[], T]) -> T: взятый ради одного SELECT'а кэша, живёт до конца функции — включая ожидание чужого HTTP (замер 11.09: 8.5 с на догрузку Яндекса), и пиковый спрос считается не в коротких SELECT'ах, а в целых фетчах. - 3. Отмена шага ДОЖИДАЕТСЯ потока (`asyncio.shield` ниже). Поток отменить - нельзя, а сессия у шага не своя — общая со всем остальным запросом. + 3. Отмена шага ДОЖИДАЕТСЯ потока — это `run_db_thread` (app/core/db.py, + #3449), общий для эстиматора и геокодера. Поток отменить нельзя, а + сессия у шага не своя — общая со всем остальным запросом. ЧТО ЭТО МЕНЯЕТ ДЛЯ ВЫЗЫВАЮЩЕГО: незакоммиченная работа его транзакции (например fill-null backfill ФИАСа в `estimate_quality`) фиксируется здесь, @@ -824,30 +825,7 @@ async def _db_step[T](db: Session, work: Callable[[], T]) -> T: 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 + return await run_db_thread(_run) async def _get_or_fetch_imv_cached( @@ -4421,7 +4399,7 @@ async def estimate_quality( if geo is None: # Без координат не можем искать через PostGIS. Возвращаем low confidence. logger.warning("geocode failed for %s — returning low-confidence estimate", payload.address) - return await asyncio.to_thread( + return await run_db_thread( _empty_estimate, payload, db, @@ -4452,7 +4430,7 @@ async def estimate_quality( dadata_fias = dadata.house_fias_id if dadata else None match_fias = payload.target_fias_id or dadata_fias try: - match = await asyncio.to_thread( + match = await run_db_thread( match_house_readonly, db, address=(dadata.canonical_address if dadata else None) or geo.full_address, @@ -4493,7 +4471,7 @@ async def estimate_quality( ) try: - await asyncio.to_thread(_backfill_house_fias) + await run_db_thread(_backfill_house_fias) except Exception as exc: # pragma: no cover — defensive logger.warning("estimate fias back-fill failed (graceful): %s", exc) except Exception as exc: # pragma: no cover — defensive @@ -4513,7 +4491,7 @@ async def estimate_quality( # house_metadata-кэше). payload остаётся главнее в любом случае — houses # заполняет только то, что пользователь не указал. if target_year is None or target_house_type is None or target_total_floors is None: - house_facts = await asyncio.to_thread( + house_facts = await run_db_thread( _lookup_house_facts, db, target_house_id=target_house_id, @@ -4585,7 +4563,7 @@ async def estimate_quality( if cohort_range is not None: cy_min, cy_max = cohort_range - listings_tier0, _, analog_tier = await asyncio.to_thread( + listings_tier0, _, analog_tier = await run_db_thread( _fetch_analogs, db, lat=geo.lat, @@ -4610,7 +4588,7 @@ async def estimate_quality( fallback_used = False else: # Tier 0 пуст/мал — graceful fallback на Tier A без cohort - listings, fallback_used, analog_tier = await asyncio.to_thread( + listings, fallback_used, analog_tier = await run_db_thread( _fetch_analogs, db, lat=geo.lat, @@ -4630,7 +4608,7 @@ async def estimate_quality( area_widened = False if len(listings) < 5: - listings_wide, _, analog_tier_wide = await asyncio.to_thread( + listings_wide, _, analog_tier_wide = await run_db_thread( _fetch_analogs, db, lat=geo.lat, @@ -4653,7 +4631,7 @@ async def estimate_quality( # Tier C: если даже на 2км мало — расширяем area tolerance до ±25% # (актуально для отдалённых районов / новостроек с нестандартной планировкой) if len(listings) < 3: - listings_widearea, _, analog_tier_wa = await asyncio.to_thread( + listings_widearea, _, analog_tier_wa = await run_db_thread( _fetch_analogs, db, lat=geo.lat, @@ -4708,7 +4686,7 @@ async def estimate_quality( """Один шаг каскада #oblast-F. Возвращает (listings, tier) только если кандидат СТРОГО больше текущей выборки — иначе релаксация не засчитана (ничего реально не выиграла) и вызывающий её не применяет.""" - candidate, _, tier = await asyncio.to_thread( + candidate, _, tier = await run_db_thread( _fetch_analogs, db, lat=geo.lat, @@ -4838,7 +4816,7 @@ async def estimate_quality( # ничего), т.е. держим гарантию «не резолвится → 66», не пусто. target_region = regions_mod.region_for_point(geo.lat, geo.lon) target_region_code = target_region.code if target_region else regions_mod.DEFAULT_REGION_CODE - dkp_raw = await asyncio.to_thread( + dkp_raw = await run_db_thread( _fetch_dkp_corridor, db, address=payload.address, @@ -4909,7 +4887,7 @@ async def estimate_quality( ), ) if yandex_val is not None: - saved_hist = await asyncio.to_thread(_save_yandex_history_items, db, yandex_val) + saved_hist = await run_db_thread(_save_yandex_history_items, db, yandex_val) logger.info( "yandex_valuation: history items processed=%d saved=%d (ext_val house_id=%s)", len(yandex_val.history_items), @@ -4992,7 +4970,7 @@ async def estimate_quality( _anchor_comps_pre: list[dict] _anchor_tier_pre: str | None if payload.area_m2 and (payload.address or (geo is not None and geo.full_address)): - _anchor_comps_pre, _anchor_tier_pre = await asyncio.to_thread( + _anchor_comps_pre, _anchor_tier_pre = await run_db_thread( _fetch_anchor_comps, db, address=payload.address or geo.full_address, @@ -5006,7 +4984,7 @@ async def estimate_quality( _anchor_comps_pre, _anchor_tier_pre = [], None # ── Pre-fetch: house IMV anchor ONCE (used for both blend and display) ─── - imv_anchor_data = await asyncio.to_thread( + imv_anchor_data = await run_db_thread( _fetch_house_imv_anchor, db, target_house_id=target_house_id, @@ -5029,7 +5007,7 @@ async def estimate_quality( and not target_quarter_cadnum and geo is not None ): - target_quarter_cadnum = await asyncio.to_thread( + target_quarter_cadnum = await run_db_thread( _lookup_target_quarter_by_coords, db, geo.lat, geo.lon ) @@ -5055,7 +5033,7 @@ async def estimate_quality( ) # ── Deterministic pricing orchestration ────────────────────────────────── - pr = await asyncio.to_thread( + pr = await run_db_thread( _price_from_inputs, listings=listings, area_m2=payload.area_m2, @@ -5175,7 +5153,7 @@ async def estimate_quality( # 5. Deals — ДКП-only sales (вторичка) из rosreestr_deals. # Importer фильтрует doc_type='ДКП' (PR-A 2026-05-24), ДДУ застройщиков # исключены — больше не скёюят median вторички ~110-120 К/м². - deals = await asyncio.to_thread( + deals = await run_db_thread( _fetch_deals, db, lat=geo.lat, @@ -5381,7 +5359,7 @@ async def estimate_quality( db.commit() - await asyncio.to_thread(_persist_estimate_and_commit) + await run_db_thread(_persist_estimate_and_commit) logger.info( "estimate: id=%s addr=%s rooms=%d area=%.1f → median=%d (n=%d, conf=%s)%s%s", @@ -5403,7 +5381,7 @@ async def estimate_quality( last_scraped_at = _compute_last_scraped_at(metadata_lots) # Месячный ₽/м² тренд целевого дома (web TREND chart) — best-effort, None если нет данных. # #audit-3: передаём freshness_months из настроек — исключаем устаревшие items. - price_trend_raw = await asyncio.to_thread( + price_trend_raw = await run_db_thread( _fetch_price_trend, db, target_house_id=target_house_id, @@ -5425,7 +5403,7 @@ async def estimate_quality( premium_building, premium_building_median_ppm2, premium_building_class, - ) = await asyncio.to_thread(_is_premium_building, db, target_house_id) + ) = await run_db_thread(_is_premium_building, db, target_house_id) # #audit-2: структурный analog_tier — стабильный enum для фронта. # anchor-путь: anchor_tier "A" → "same_building", "C" → "micro_radius". diff --git a/tradein-mvp/backend/app/services/geocoder.py b/tradein-mvp/backend/app/services/geocoder.py index 6d29ff8b..f78b0da4 100644 --- a/tradein-mvp/backend/app/services/geocoder.py +++ b/tradein-mvp/backend/app/services/geocoder.py @@ -26,6 +26,7 @@ from sqlalchemy.orm import Session from tenacity import retry, stop_after_attempt, wait_exponential from app.core.config import settings +from app.core.db import run_db_thread from app.services import dadata from app.services.regions import REGIONS as _ALL_REGIONS from app.services.regions import Region, is_within_bbox @@ -1884,11 +1885,11 @@ async def suggest( parsed = _parse_street_house(query.strip()) if parsed is not None: street, house = parsed - hit = await asyncio.to_thread(_cadastral_house_match, db, street, house) + hit = await run_db_thread(_cadastral_house_match, db, street, house) if hit is not None: return [hit] # 1b. Fallback: legacy raw-ILIKE forward search (для нераспарсенных форм) - cad_results = await asyncio.to_thread(_cadastral_forward_sync, db, query.strip(), limit) + cad_results = await run_db_thread(_cadastral_forward_sync, db, query.strip(), limit) if cad_results: return cad_results @@ -1995,7 +1996,7 @@ async def _geocode_resolve( addr_norm = _cache_key(normalize_address(address), city_hint) # 1. Cache (sync DB-IO → offload в threadpool, чтобы не блокировать event loop) - cached = await asyncio.to_thread(_cache_get, db, addr_norm) + cached = await run_db_thread(_cache_get, db, addr_norm) if cached is not None: logger.info("geocode cache hit: %s", addr_norm) return replace(cached, city_ambiguous=city_ambiguous) @@ -2028,7 +2029,7 @@ async def _geocode_resolve( if use_local_ekb and parsed is not None: street, house = parsed try: - hit = await asyncio.to_thread(_geoportal_house_match, db, street, house) + hit = await run_db_thread(_geoportal_house_match, db, street, house) except Exception: logger.warning("geoportal house-match raised — fall through", exc_info=True) hit = None @@ -2041,7 +2042,7 @@ async def _geocode_resolve( confidence="exact", city_ambiguous=city_ambiguous, ) - await asyncio.to_thread(_cache_put, db, addr_norm, result) + await run_db_thread(_cache_put, db, addr_norm, result) logger.info( "geocode geoportal house-match: %s → (%.5f, %.5f)", addr_norm, @@ -2056,7 +2057,7 @@ async def _geocode_resolve( # (литеральная подстрока не совпадает). if use_local_ekb and parsed is not None: street, house = parsed - hit = await asyncio.to_thread(_cadastral_house_match, db, street, house) + hit = await run_db_thread(_cadastral_house_match, db, street, house) if hit is not None: result = GeocodeResult( lat=hit.lat, @@ -2066,7 +2067,7 @@ async def _geocode_resolve( confidence="exact", city_ambiguous=city_ambiguous, ) - await asyncio.to_thread(_cache_put, db, addr_norm, result) + await run_db_thread(_cache_put, db, addr_norm, result) logger.info( "geocode cadastral house-match: %s → (%.5f, %.5f)", addr_norm, @@ -2077,9 +2078,7 @@ async def _geocode_resolve( # 2d. Fallback: legacy raw-ILIKE forward search (для нераспарсенных форм) if use_local_ekb: - cad_suggestions = await asyncio.to_thread( - _cadastral_forward_sync, db, address.strip(), limit=1 - ) + cad_suggestions = await run_db_thread(_cadastral_forward_sync, db, address.strip(), limit=1) if cad_suggestions: s = cad_suggestions[0] result = GeocodeResult( @@ -2090,7 +2089,7 @@ async def _geocode_resolve( confidence="exact", city_ambiguous=city_ambiguous, ) - await asyncio.to_thread(_cache_put, db, addr_norm, result) + await run_db_thread(_cache_put, db, addr_norm, result) logger.info( "geocode cadastral fdw: %s → (%.5f, %.5f)", addr_norm, result.lat, result.lon ) @@ -2101,7 +2100,7 @@ async def _geocode_resolve( result = await _nominatim_lookup(address, city_hint, region_code) if result is not None: result = replace(result, city_ambiguous=city_ambiguous) - await asyncio.to_thread(_cache_put, db, addr_norm, result) + await run_db_thread(_cache_put, db, addr_norm, result) logger.info("geocode nominatim: %s → (%.5f, %.5f)", addr_norm, result.lat, result.lon) return result except Exception: @@ -2120,7 +2119,7 @@ async def _geocode_resolve( if use_local_ekb and parsed is not None: local_street, _parsed_house = parsed local_house = _extract_local_house_token(address) or _parsed_house - hit = await asyncio.to_thread(_local_houses_match, db, local_street, local_house) + hit = await run_db_thread(_local_houses_match, db, local_street, local_house) if hit is not None: result = GeocodeResult( lat=hit.lat, @@ -2329,7 +2328,7 @@ async def reverse_geocode( """ # 1. Cadastral FDW primary (без внешнего API, возвращает жилой дом not POI) if db is not None: - cad = await asyncio.to_thread(_cadastral_reverse_sync_full, db, lat, lon) + cad = await run_db_thread(_cadastral_reverse_sync_full, db, lat, lon) if cad is not None: address, snap_lat, snap_lon = cad return ReverseGeocodeResult( diff --git a/tradein-mvp/backend/tests/test_3449_geocoder_cancel_orphan.py b/tradein-mvp/backend/tests/test_3449_geocoder_cancel_orphan.py new file mode 100644 index 00000000..90e380f4 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3449_geocoder_cancel_orphan.py @@ -0,0 +1,88 @@ +"""Отмена геокодинга по бюджету не оставляет ОСИРОТЕВШИЙ поток в сессии запроса (#3449). + +`geocode()` живёт под бюджетом 12 с (`estimator._with_budget` = `asyncio.wait_for`), а +внутри ходит в БД (`_cache_get`/`_cache_put`/локальные тиры) по ТОЙ ЖЕ `Session`, с +которой запрос идёт дальше. `asyncio.to_thread` отменить нельзя: по истечении бюджета +снимается только ожидание со стороны loop'а, поток продолжает работать в чужой сессии. +Ошибка геокодера при этом никого не трогает (её глушит `_with_budget`) — страдает +СЛЕДУЮЩИЙ потребитель сессии, и у `_persist_estimate_and_commit` её не ловит никто: +500 и потерянная оценка клиента. + +Меряем значение, а не форму: «следующий шаг не вошёл в сессию, пока сирота не +закончил». Следующий шаг здесь — ГОЛЫЙ `asyncio.to_thread(db...)`, как персист оценки, +а не ещё один защищённый вызов: защита, которая живёт только внутри обёртки, ровно +того пострадавшего и не закрывает. +""" + +from __future__ import annotations + +import asyncio +import os +import threading +import time +from typing import Any + +import pytest + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +from app.services import estimator as est +from app.services import geocoder as geo + +# Шаг БД заметно длиннее бюджета — окно, в котором сирота ещё работает. +_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 + + 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 + + +async def test_geocode_budget_cancel_does_not_leave_orphan_in_session( + monkeypatch: pytest.MonkeyPatch, +) -> None: + db = _ConcurrencyProbeSession() + + # Настоящий путь `geocode()` до первого шага БД; сам SQL подменён — меряем + # владение сессией, а не содержимое кэша. + def _slow_cache_get(session: Any, address_norm: str) -> None: + session.execute("geocode cache lookup", {"addr": address_norm}) + return None + + monkeypatch.setattr(geo, "_cache_get", _slow_cache_get) + + degraded = await est._with_budget( + geo.geocode("ул. Пушкина, д. 1", db), # type: ignore[arg-type] + _BUDGET_S, + label="geocode", + ) + # Бюджет истёк — геокодер деградировал в None, вызывающий идёт дальше. + assert degraded is None, "бюджет не сработал — тест ничего не проверил" + + # Следующий шаг ТОГО ЖЕ запроса по ТОЙ ЖЕ сессии (образец — персист оценки). + 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» на персисте" + ) From be9aa2f9071e59a1b234bcb765588a8651d5e4ee Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 12 Sep 2026 13:57:49 +0500 Subject: [PATCH 2/2] =?UTF-8?q?=D0=93=D0=B5=D0=B9=D1=82=20=D0=BD=D0=B0=20?= =?UTF-8?q?=D0=92=D0=A1=D0=95=2034=20=D0=BF=D1=80=D0=BE=D0=B2=D0=BE=D0=B4?= =?UTF-8?q?=D0=BA=D0=B8=20+=20=D0=B7=D0=B0=D0=BF=D1=80=D0=B5=D1=82=20?= =?UTF-8?q?=D0=B2=D0=BB=D0=BE=D0=B6=D0=B5=D0=BD=D0=BD=D1=8B=D1=85=20=D0=B1?= =?UTF-8?q?=D1=8E=D0=B4=D0=B6=D0=B5=D1=82=D0=BE=D0=B2=20(=D1=80=D0=B5?= =?UTF-8?q?=D0=B2=D1=8C=D1=8E=20#3460)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Сценарный тест ловил одну проводку из 34 — ту, через которую сам и шёл (`geocoder._cache_get`). Мутационный прогон ревьюера: возврат голого `asyncio.to_thread` в 5 из 6 других мест тест НЕ краснит, то есть регресс «кто-то вернул вызов в голый вид» прошёл бы мимо CI в 33 случаях из 34. `test_no_bare_to_thread_over_request_session` читает исходники geocoder и estimator (через `module.__file__`, не по относительному пути — он зависел бы от cwd прогона) и требует нуля живых `asyncio.to_thread(`. Оба модуля сейчас на нуле, поэтому гейт без списка исключений. Фальсификация — голый `to_thread` у `_fetch_anchor_comps` (estimator:4973, сценарным тестом не покрыт): гейт краснеет с номером строки. Второе: защита `run_db_thread` одноразовая — `except asyncio.CancelledError` ловит ОДНУ отмену, вторая вылетает из самого `asyncio.wait([step])`, и поток остаётся сиротой. Живых путей нет (`_with_budget` нигде не вложен, Starlette не отменяет задачу на дисконнекте, uvicorn стартует без `--timeout-graceful-shutdown`), поэтому кода не трогаю — фиксирую инвариант «не вкладывать бюджеты» в докстринге `_with_budget`, чтобы вложение не завезли как безобидное. Co-Authored-By: Claude Opus 5 --- tradein-mvp/backend/app/services/estimator.py | 6 +++++ .../tests/test_3449_geocoder_cancel_orphan.py | 27 +++++++++++++++++++ 2 files changed, 33 insertions(+) diff --git a/tradein-mvp/backend/app/services/estimator.py b/tradein-mvp/backend/app/services/estimator.py index 87426284..c20fc154 100644 --- a/tradein-mvp/backend/app/services/estimator.py +++ b/tradein-mvp/backend/app/services/estimator.py @@ -3092,6 +3092,12 @@ async def _with_budget(coro: Any, budget_s: float, *, label: str) -> Any: blowing the gateway read timeout (#654: opaque Caddy 502/504). budget_s <= 0 disables the guard (await directly) — escape hatch via config. + + НЕ ВКЛАДЫВАТЬ бюджеты друг в друга: защита `run_db_thread` (#3449) одноразовая + — она ловит ОДНУ отмену, а вторая, прилетевшая пока шаг БД дожидается своего + потока, вылетает из самого ожидания, и поток остаётся сиротой в чужой `Session` + (ровно то, ради чего защита и заведена). Сейчас ни один вызов `_with_budget` не + обёрнут другим — это инвариант, а не совпадение. """ if budget_s is None or budget_s <= 0: return await coro diff --git a/tradein-mvp/backend/tests/test_3449_geocoder_cancel_orphan.py b/tradein-mvp/backend/tests/test_3449_geocoder_cancel_orphan.py index 90e380f4..3f022203 100644 --- a/tradein-mvp/backend/tests/test_3449_geocoder_cancel_orphan.py +++ b/tradein-mvp/backend/tests/test_3449_geocoder_cancel_orphan.py @@ -18,6 +18,7 @@ from __future__ import annotations import asyncio import os +import pathlib import threading import time from typing import Any @@ -86,3 +87,29 @@ async def test_geocode_budget_cancel_does_not_leave_orphan_in_session( "персист вошёл в сессию, пока осиротевший поток геокодера ещё работал в ней: " "два потока в одной Session → «another operation is in progress» на персисте" ) + + +def test_no_bare_to_thread_over_request_session() -> None: + """Source-гейт: в геокодере и эстиматоре не осталось голых `asyncio.to_thread(`. + + Тест выше ловит ОДНУ проводку — ту, через которую идёт сценарий. Остальные 33 + (`_cache_put`, `_fetch_anchor_comps`, персист, …) он не видит: возврат любой из + них в голый вид прошёл бы мимо CI. Оба модуля сейчас на нуле по живым вызовам, + поэтому гейт — ровно «ноль», без списка исключений. Понадобится вызов со СВОЕЙ + сессией (как `user_events.record_event`) — заводить его в отдельном модуле или + менять этот тест осознанно. + + Читаем через `module.__file__`: относительный путь зависел бы от cwd прогона. + """ + for module in (geo, est): + src = pathlib.Path(module.__file__ or "").read_text(encoding="utf-8") + bare = [ + f"{i}: {line.strip()}" + for i, line in enumerate(src.splitlines(), 1) + if "asyncio.to_thread(" in line and not line.lstrip().startswith("#") + ] + assert not bare, ( + f"{module.__name__}: голый asyncio.to_thread по сессии запроса — " + f"отмена оставит сироту в чужой Session (#3449), нужен run_db_thread:\n" + + "\n".join(bare) + )