Оценка не теряется, когда геокодер не уложился в бюджет (#3449) #3460
6 changed files with 205 additions and 68 deletions
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
@ -3114,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
|
||||
|
|
@ -4421,7 +4405,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 +4436,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 +4477,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 +4497,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 +4569,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 +4594,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 +4614,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 +4637,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 +4692,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 +4822,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 +4893,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 +4976,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 +4990,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 +5013,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 +5039,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 +5159,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 +5365,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 +5387,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 +5409,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".
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
115
tradein-mvp/backend/tests/test_3449_geocoder_cancel_orphan.py
Normal file
115
tradein-mvp/backend/tests/test_3449_geocoder_cancel_orphan.py
Normal file
|
|
@ -0,0 +1,115 @@
|
|||
"""Отмена геокодинга по бюджету не оставляет ОСИРОТЕВШИЙ поток в сессии запроса (#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 pathlib
|
||||
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» на персисте"
|
||||
)
|
||||
|
||||
|
||||
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)
|
||||
)
|
||||
Loading…
Add table
Reference in a new issue