Оценка не теряется, когда геокодер не уложился в бюджет (#3449) #3460
6 changed files with 205 additions and 68 deletions
|
|
@ -2,7 +2,6 @@
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import asyncio
|
|
||||||
import logging
|
import logging
|
||||||
from typing import Annotated
|
from typing import Annotated
|
||||||
|
|
||||||
|
|
@ -11,7 +10,7 @@ from pydantic import BaseModel, Field
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
from sqlalchemy.orm import Session
|
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.estimator import _lookup_house_facts
|
||||||
from app.services.geocoder import GeocodeResult, geocode, reverse_geocode, suggest
|
from app.services.geocoder import GeocodeResult, geocode, reverse_geocode, suggest
|
||||||
|
|
||||||
|
|
@ -254,9 +253,9 @@ async def house_facts(
|
||||||
"""
|
"""
|
||||||
target_house_id: int | None = None
|
target_house_id: int | None = None
|
||||||
if fias_id is not 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,
|
_lookup_house_facts,
|
||||||
db,
|
db,
|
||||||
target_house_id=target_house_id,
|
target_house_id=target_house_id,
|
||||||
|
|
|
||||||
|
|
@ -14,7 +14,6 @@ module docstring (честно про то, что не всегда разре
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import asyncio
|
|
||||||
import logging
|
import logging
|
||||||
from typing import Annotated
|
from typing import Annotated
|
||||||
from uuid import UUID
|
from uuid import UUID
|
||||||
|
|
@ -23,7 +22,7 @@ from fastapi import APIRouter, Depends, HTTPException
|
||||||
from pydantic import BaseModel, Field
|
from pydantic import BaseModel, Field
|
||||||
from sqlalchemy.orm import Session
|
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__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
@ -67,7 +66,7 @@ async def erase_person_data_endpoint(
|
||||||
|
|
||||||
from app.services.data_erasure import erase_person_data
|
from app.services.data_erasure import erase_person_data
|
||||||
|
|
||||||
counters = await asyncio.to_thread(
|
counters = await run_db_thread(
|
||||||
erase_person_data,
|
erase_person_data,
|
||||||
db,
|
db,
|
||||||
username=payload.username,
|
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 import create_engine
|
||||||
from sqlalchemy.orm import DeclarativeBase, Session, sessionmaker
|
from sqlalchemy.orm import DeclarativeBase, Session, sessionmaker
|
||||||
|
|
||||||
from app.core.config import settings
|
from app.core.config import settings
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
engine = create_engine(
|
engine = create_engine(
|
||||||
settings.database_url,
|
settings.database_url,
|
||||||
pool_pre_ping=True,
|
pool_pre_ping=True,
|
||||||
|
|
@ -65,3 +70,39 @@ def get_db() -> Generator[Session, None, None]:
|
||||||
yield db
|
yield db
|
||||||
finally:
|
finally:
|
||||||
db.close()
|
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 sqlalchemy.orm import Session
|
||||||
|
|
||||||
from app.core.config import LISTINGS_FRESH_DAYS, Settings, settings
|
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 (
|
from app.schemas.trade_in import (
|
||||||
AggregatedEstimate,
|
AggregatedEstimate,
|
||||||
AnalogLot,
|
AnalogLot,
|
||||||
|
|
@ -803,8 +803,9 @@ 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` ниже). Поток отменить
|
3. Отмена шага ДОЖИДАЕТСЯ потока — это `run_db_thread` (app/core/db.py,
|
||||||
нельзя, а сессия у шага не своя — общая со всем остальным запросом.
|
#3449), общий для эстиматора и геокодера. Поток отменить нельзя, а
|
||||||
|
сессия у шага не своя — общая со всем остальным запросом.
|
||||||
|
|
||||||
ЧТО ЭТО МЕНЯЕТ ДЛЯ ВЫЗЫВАЮЩЕГО: незакоммиченная работа его транзакции
|
ЧТО ЭТО МЕНЯЕТ ДЛЯ ВЫЗЫВАЮЩЕГО: незакоммиченная работа его транзакции
|
||||||
(например fill-null backfill ФИАСа в `estimate_quality`) фиксируется здесь,
|
(например fill-null backfill ФИАСа в `estimate_quality`) фиксируется здесь,
|
||||||
|
|
@ -824,30 +825,7 @@ async def _db_step[T](db: Session, work: Callable[[], T]) -> T:
|
||||||
db.commit()
|
db.commit()
|
||||||
return result
|
return result
|
||||||
|
|
||||||
step = asyncio.ensure_future(asyncio.to_thread(_run))
|
return await run_db_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(
|
||||||
|
|
@ -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).
|
blowing the gateway read timeout (#654: opaque Caddy 502/504).
|
||||||
|
|
||||||
budget_s <= 0 disables the guard (await directly) — escape hatch via config.
|
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:
|
if budget_s is None or budget_s <= 0:
|
||||||
return await coro
|
return await coro
|
||||||
|
|
@ -4421,7 +4405,7 @@ async def estimate_quality(
|
||||||
if geo is None:
|
if geo is None:
|
||||||
# Без координат не можем искать через PostGIS. Возвращаем low confidence.
|
# Без координат не можем искать через PostGIS. Возвращаем low confidence.
|
||||||
logger.warning("geocode failed for %s — returning low-confidence estimate", payload.address)
|
logger.warning("geocode failed for %s — returning low-confidence estimate", payload.address)
|
||||||
return await asyncio.to_thread(
|
return await run_db_thread(
|
||||||
_empty_estimate,
|
_empty_estimate,
|
||||||
payload,
|
payload,
|
||||||
db,
|
db,
|
||||||
|
|
@ -4452,7 +4436,7 @@ async def estimate_quality(
|
||||||
dadata_fias = dadata.house_fias_id if dadata else None
|
dadata_fias = dadata.house_fias_id if dadata else None
|
||||||
match_fias = payload.target_fias_id or dadata_fias
|
match_fias = payload.target_fias_id or dadata_fias
|
||||||
try:
|
try:
|
||||||
match = await asyncio.to_thread(
|
match = await run_db_thread(
|
||||||
match_house_readonly,
|
match_house_readonly,
|
||||||
db,
|
db,
|
||||||
address=(dadata.canonical_address if dadata else None) or geo.full_address,
|
address=(dadata.canonical_address if dadata else None) or geo.full_address,
|
||||||
|
|
@ -4493,7 +4477,7 @@ async def estimate_quality(
|
||||||
)
|
)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
await asyncio.to_thread(_backfill_house_fias)
|
await run_db_thread(_backfill_house_fias)
|
||||||
except Exception as exc: # pragma: no cover — defensive
|
except Exception as exc: # pragma: no cover — defensive
|
||||||
logger.warning("estimate fias back-fill failed (graceful): %s", exc)
|
logger.warning("estimate fias back-fill failed (graceful): %s", exc)
|
||||||
except Exception as exc: # pragma: no cover — defensive
|
except Exception as exc: # pragma: no cover — defensive
|
||||||
|
|
@ -4513,7 +4497,7 @@ async def estimate_quality(
|
||||||
# house_metadata-кэше). payload остаётся главнее в любом случае — houses
|
# house_metadata-кэше). payload остаётся главнее в любом случае — houses
|
||||||
# заполняет только то, что пользователь не указал.
|
# заполняет только то, что пользователь не указал.
|
||||||
if target_year is None or target_house_type is None or target_total_floors is None:
|
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,
|
_lookup_house_facts,
|
||||||
db,
|
db,
|
||||||
target_house_id=target_house_id,
|
target_house_id=target_house_id,
|
||||||
|
|
@ -4585,7 +4569,7 @@ async def estimate_quality(
|
||||||
|
|
||||||
if cohort_range is not None:
|
if cohort_range is not None:
|
||||||
cy_min, cy_max = cohort_range
|
cy_min, cy_max = cohort_range
|
||||||
listings_tier0, _, analog_tier = await asyncio.to_thread(
|
listings_tier0, _, analog_tier = await run_db_thread(
|
||||||
_fetch_analogs,
|
_fetch_analogs,
|
||||||
db,
|
db,
|
||||||
lat=geo.lat,
|
lat=geo.lat,
|
||||||
|
|
@ -4610,7 +4594,7 @@ async def estimate_quality(
|
||||||
fallback_used = False
|
fallback_used = False
|
||||||
else:
|
else:
|
||||||
# Tier 0 пуст/мал — graceful fallback на Tier A без cohort
|
# 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,
|
_fetch_analogs,
|
||||||
db,
|
db,
|
||||||
lat=geo.lat,
|
lat=geo.lat,
|
||||||
|
|
@ -4630,7 +4614,7 @@ async def estimate_quality(
|
||||||
area_widened = False
|
area_widened = False
|
||||||
|
|
||||||
if len(listings) < 5:
|
if len(listings) < 5:
|
||||||
listings_wide, _, analog_tier_wide = await asyncio.to_thread(
|
listings_wide, _, analog_tier_wide = await run_db_thread(
|
||||||
_fetch_analogs,
|
_fetch_analogs,
|
||||||
db,
|
db,
|
||||||
lat=geo.lat,
|
lat=geo.lat,
|
||||||
|
|
@ -4653,7 +4637,7 @@ async def estimate_quality(
|
||||||
# Tier C: если даже на 2км мало — расширяем area tolerance до ±25%
|
# Tier C: если даже на 2км мало — расширяем area tolerance до ±25%
|
||||||
# (актуально для отдалённых районов / новостроек с нестандартной планировкой)
|
# (актуально для отдалённых районов / новостроек с нестандартной планировкой)
|
||||||
if len(listings) < 3:
|
if len(listings) < 3:
|
||||||
listings_widearea, _, analog_tier_wa = await asyncio.to_thread(
|
listings_widearea, _, analog_tier_wa = await run_db_thread(
|
||||||
_fetch_analogs,
|
_fetch_analogs,
|
||||||
db,
|
db,
|
||||||
lat=geo.lat,
|
lat=geo.lat,
|
||||||
|
|
@ -4708,7 +4692,7 @@ async def estimate_quality(
|
||||||
"""Один шаг каскада #oblast-F. Возвращает (listings, tier) только если
|
"""Один шаг каскада #oblast-F. Возвращает (listings, tier) только если
|
||||||
кандидат СТРОГО больше текущей выборки — иначе релаксация не засчитана
|
кандидат СТРОГО больше текущей выборки — иначе релаксация не засчитана
|
||||||
(ничего реально не выиграла) и вызывающий её не применяет."""
|
(ничего реально не выиграла) и вызывающий её не применяет."""
|
||||||
candidate, _, tier = await asyncio.to_thread(
|
candidate, _, tier = await run_db_thread(
|
||||||
_fetch_analogs,
|
_fetch_analogs,
|
||||||
db,
|
db,
|
||||||
lat=geo.lat,
|
lat=geo.lat,
|
||||||
|
|
@ -4838,7 +4822,7 @@ async def estimate_quality(
|
||||||
# ничего), т.е. держим гарантию «не резолвится → 66», не пусто.
|
# ничего), т.е. держим гарантию «не резолвится → 66», не пусто.
|
||||||
target_region = regions_mod.region_for_point(geo.lat, geo.lon)
|
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
|
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,
|
_fetch_dkp_corridor,
|
||||||
db,
|
db,
|
||||||
address=payload.address,
|
address=payload.address,
|
||||||
|
|
@ -4909,7 +4893,7 @@ async def estimate_quality(
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
if yandex_val is not None:
|
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(
|
logger.info(
|
||||||
"yandex_valuation: history items processed=%d saved=%d (ext_val house_id=%s)",
|
"yandex_valuation: history items processed=%d saved=%d (ext_val house_id=%s)",
|
||||||
len(yandex_val.history_items),
|
len(yandex_val.history_items),
|
||||||
|
|
@ -4992,7 +4976,7 @@ async def estimate_quality(
|
||||||
_anchor_comps_pre: list[dict]
|
_anchor_comps_pre: list[dict]
|
||||||
_anchor_tier_pre: str | None
|
_anchor_tier_pre: str | None
|
||||||
if payload.area_m2 and (payload.address or (geo is not None and geo.full_address)):
|
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,
|
_fetch_anchor_comps,
|
||||||
db,
|
db,
|
||||||
address=payload.address or geo.full_address,
|
address=payload.address or geo.full_address,
|
||||||
|
|
@ -5006,7 +4990,7 @@ async def estimate_quality(
|
||||||
_anchor_comps_pre, _anchor_tier_pre = [], None
|
_anchor_comps_pre, _anchor_tier_pre = [], None
|
||||||
|
|
||||||
# ── Pre-fetch: house IMV anchor ONCE (used for both blend and display) ───
|
# ── 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,
|
_fetch_house_imv_anchor,
|
||||||
db,
|
db,
|
||||||
target_house_id=target_house_id,
|
target_house_id=target_house_id,
|
||||||
|
|
@ -5029,7 +5013,7 @@ async def estimate_quality(
|
||||||
and not target_quarter_cadnum
|
and not target_quarter_cadnum
|
||||||
and geo is not None
|
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
|
_lookup_target_quarter_by_coords, db, geo.lat, geo.lon
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -5055,7 +5039,7 @@ async def estimate_quality(
|
||||||
)
|
)
|
||||||
|
|
||||||
# ── Deterministic pricing orchestration ──────────────────────────────────
|
# ── Deterministic pricing orchestration ──────────────────────────────────
|
||||||
pr = await asyncio.to_thread(
|
pr = await run_db_thread(
|
||||||
_price_from_inputs,
|
_price_from_inputs,
|
||||||
listings=listings,
|
listings=listings,
|
||||||
area_m2=payload.area_m2,
|
area_m2=payload.area_m2,
|
||||||
|
|
@ -5175,7 +5159,7 @@ async def estimate_quality(
|
||||||
# 5. Deals — ДКП-only sales (вторичка) из rosreestr_deals.
|
# 5. Deals — ДКП-only sales (вторичка) из rosreestr_deals.
|
||||||
# Importer фильтрует doc_type='ДКП' (PR-A 2026-05-24), ДДУ застройщиков
|
# Importer фильтрует doc_type='ДКП' (PR-A 2026-05-24), ДДУ застройщиков
|
||||||
# исключены — больше не скёюят median вторички ~110-120 К/м².
|
# исключены — больше не скёюят median вторички ~110-120 К/м².
|
||||||
deals = await asyncio.to_thread(
|
deals = await run_db_thread(
|
||||||
_fetch_deals,
|
_fetch_deals,
|
||||||
db,
|
db,
|
||||||
lat=geo.lat,
|
lat=geo.lat,
|
||||||
|
|
@ -5381,7 +5365,7 @@ async def estimate_quality(
|
||||||
|
|
||||||
db.commit()
|
db.commit()
|
||||||
|
|
||||||
await asyncio.to_thread(_persist_estimate_and_commit)
|
await run_db_thread(_persist_estimate_and_commit)
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"estimate: id=%s addr=%s rooms=%d area=%.1f → median=%d (n=%d, conf=%s)%s%s",
|
"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)
|
last_scraped_at = _compute_last_scraped_at(metadata_lots)
|
||||||
# Месячный ₽/м² тренд целевого дома (web TREND chart) — best-effort, None если нет данных.
|
# Месячный ₽/м² тренд целевого дома (web TREND chart) — best-effort, None если нет данных.
|
||||||
# #audit-3: передаём freshness_months из настроек — исключаем устаревшие items.
|
# #audit-3: передаём freshness_months из настроек — исключаем устаревшие items.
|
||||||
price_trend_raw = await asyncio.to_thread(
|
price_trend_raw = await run_db_thread(
|
||||||
_fetch_price_trend,
|
_fetch_price_trend,
|
||||||
db,
|
db,
|
||||||
target_house_id=target_house_id,
|
target_house_id=target_house_id,
|
||||||
|
|
@ -5425,7 +5409,7 @@ async def estimate_quality(
|
||||||
premium_building,
|
premium_building,
|
||||||
premium_building_median_ppm2,
|
premium_building_median_ppm2,
|
||||||
premium_building_class,
|
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 для фронта.
|
# #audit-2: структурный analog_tier — стабильный enum для фронта.
|
||||||
# anchor-путь: anchor_tier "A" → "same_building", "C" → "micro_radius".
|
# 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 tenacity import retry, stop_after_attempt, wait_exponential
|
||||||
|
|
||||||
from app.core.config import settings
|
from app.core.config import settings
|
||||||
|
from app.core.db import run_db_thread
|
||||||
from app.services import dadata
|
from app.services import dadata
|
||||||
from app.services.regions import REGIONS as _ALL_REGIONS
|
from app.services.regions import REGIONS as _ALL_REGIONS
|
||||||
from app.services.regions import Region, is_within_bbox
|
from app.services.regions import Region, is_within_bbox
|
||||||
|
|
@ -1884,11 +1885,11 @@ async def suggest(
|
||||||
parsed = _parse_street_house(query.strip())
|
parsed = _parse_street_house(query.strip())
|
||||||
if parsed is not None:
|
if parsed is not None:
|
||||||
street, house = parsed
|
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:
|
if hit is not None:
|
||||||
return [hit]
|
return [hit]
|
||||||
# 1b. Fallback: legacy raw-ILIKE forward search (для нераспарсенных форм)
|
# 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:
|
if cad_results:
|
||||||
return cad_results
|
return cad_results
|
||||||
|
|
||||||
|
|
@ -1995,7 +1996,7 @@ async def _geocode_resolve(
|
||||||
addr_norm = _cache_key(normalize_address(address), city_hint)
|
addr_norm = _cache_key(normalize_address(address), city_hint)
|
||||||
|
|
||||||
# 1. Cache (sync DB-IO → offload в threadpool, чтобы не блокировать event loop)
|
# 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:
|
if cached is not None:
|
||||||
logger.info("geocode cache hit: %s", addr_norm)
|
logger.info("geocode cache hit: %s", addr_norm)
|
||||||
return replace(cached, city_ambiguous=city_ambiguous)
|
return replace(cached, city_ambiguous=city_ambiguous)
|
||||||
|
|
@ -2028,7 +2029,7 @@ async def _geocode_resolve(
|
||||||
if use_local_ekb and parsed is not None:
|
if use_local_ekb and parsed is not None:
|
||||||
street, house = parsed
|
street, house = parsed
|
||||||
try:
|
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:
|
except Exception:
|
||||||
logger.warning("geoportal house-match raised — fall through", exc_info=True)
|
logger.warning("geoportal house-match raised — fall through", exc_info=True)
|
||||||
hit = None
|
hit = None
|
||||||
|
|
@ -2041,7 +2042,7 @@ async def _geocode_resolve(
|
||||||
confidence="exact",
|
confidence="exact",
|
||||||
city_ambiguous=city_ambiguous,
|
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(
|
logger.info(
|
||||||
"geocode geoportal house-match: %s → (%.5f, %.5f)",
|
"geocode geoportal house-match: %s → (%.5f, %.5f)",
|
||||||
addr_norm,
|
addr_norm,
|
||||||
|
|
@ -2056,7 +2057,7 @@ async def _geocode_resolve(
|
||||||
# (литеральная подстрока не совпадает).
|
# (литеральная подстрока не совпадает).
|
||||||
if use_local_ekb and parsed is not None:
|
if use_local_ekb and parsed is not None:
|
||||||
street, house = parsed
|
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:
|
if hit is not None:
|
||||||
result = GeocodeResult(
|
result = GeocodeResult(
|
||||||
lat=hit.lat,
|
lat=hit.lat,
|
||||||
|
|
@ -2066,7 +2067,7 @@ async def _geocode_resolve(
|
||||||
confidence="exact",
|
confidence="exact",
|
||||||
city_ambiguous=city_ambiguous,
|
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(
|
logger.info(
|
||||||
"geocode cadastral house-match: %s → (%.5f, %.5f)",
|
"geocode cadastral house-match: %s → (%.5f, %.5f)",
|
||||||
addr_norm,
|
addr_norm,
|
||||||
|
|
@ -2077,9 +2078,7 @@ async def _geocode_resolve(
|
||||||
|
|
||||||
# 2d. Fallback: legacy raw-ILIKE forward search (для нераспарсенных форм)
|
# 2d. Fallback: legacy raw-ILIKE forward search (для нераспарсенных форм)
|
||||||
if use_local_ekb:
|
if use_local_ekb:
|
||||||
cad_suggestions = await asyncio.to_thread(
|
cad_suggestions = await run_db_thread(_cadastral_forward_sync, db, address.strip(), limit=1)
|
||||||
_cadastral_forward_sync, db, address.strip(), limit=1
|
|
||||||
)
|
|
||||||
if cad_suggestions:
|
if cad_suggestions:
|
||||||
s = cad_suggestions[0]
|
s = cad_suggestions[0]
|
||||||
result = GeocodeResult(
|
result = GeocodeResult(
|
||||||
|
|
@ -2090,7 +2089,7 @@ async def _geocode_resolve(
|
||||||
confidence="exact",
|
confidence="exact",
|
||||||
city_ambiguous=city_ambiguous,
|
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(
|
logger.info(
|
||||||
"geocode cadastral fdw: %s → (%.5f, %.5f)", addr_norm, result.lat, result.lon
|
"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)
|
result = await _nominatim_lookup(address, city_hint, region_code)
|
||||||
if result is not None:
|
if result is not None:
|
||||||
result = replace(result, city_ambiguous=city_ambiguous)
|
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)
|
logger.info("geocode nominatim: %s → (%.5f, %.5f)", addr_norm, result.lat, result.lon)
|
||||||
return result
|
return result
|
||||||
except Exception:
|
except Exception:
|
||||||
|
|
@ -2120,7 +2119,7 @@ async def _geocode_resolve(
|
||||||
if use_local_ekb and parsed is not None:
|
if use_local_ekb and parsed is not None:
|
||||||
local_street, _parsed_house = parsed
|
local_street, _parsed_house = parsed
|
||||||
local_house = _extract_local_house_token(address) or _parsed_house
|
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:
|
if hit is not None:
|
||||||
result = GeocodeResult(
|
result = GeocodeResult(
|
||||||
lat=hit.lat,
|
lat=hit.lat,
|
||||||
|
|
@ -2329,7 +2328,7 @@ async def reverse_geocode(
|
||||||
"""
|
"""
|
||||||
# 1. Cadastral FDW primary (без внешнего API, возвращает жилой дом not POI)
|
# 1. Cadastral FDW primary (без внешнего API, возвращает жилой дом not POI)
|
||||||
if db is not None:
|
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:
|
if cad is not None:
|
||||||
address, snap_lat, snap_lon = cad
|
address, snap_lat, snap_lon = cad
|
||||||
return ReverseGeocodeResult(
|
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