diff --git a/tradein-mvp/backend/app/core/config.py b/tradein-mvp/backend/app/core/config.py index 2c035f89..7511f325 100644 --- a/tradein-mvp/backend/app/core/config.py +++ b/tradein-mvp/backend/app/core/config.py @@ -723,6 +723,26 @@ class Settings(BaseSettings): # ESTIMATE_HOUSE_META_TIMEOUT_S. estimate_yandex_valuation_timeout_s: float = 8.0 estimate_cian_valuation_timeout_s: float = 8.0 + + # Внешние оценки (Yandex/Cian) не ждать в запросе, а догружать в фоне. + # + # Замер на проде 2026-08-22: расчёт по НОВОМУ адресу занимает 8-12 с, из них + # ~6 с ждёт Yandex и ~1.5 с Cian. Собственные запросы к базе и сам расчёт + # укладываются в секунду. По уже виденному адресу (кэш 24 ч) — 0.4-0.8 с. + # + # True: в запросе делается только чтение кэша; при промахе источник + # деградирует в None, а свежая загрузка уходит в фон и наполняет кэш к + # следующему обращению по тому же адресу. Ответ отдаётся за ~1 с. + # + # Деградация в None — НЕ новое состояние ответа: ровно так же ведёт себя + # таймаут `estimate_*_valuation_timeout_s`, этот путь работает в проде + # сегодня. Поэтому переключение не меняет контракт API. + # + # False (дефолт) — прежнее поведение: ждать источники в запросе. + # ENV: ESTIMATE_EXTERNAL_SOURCES_BACKGROUND. + estimate_external_sources_background: bool = Field( + default=False, validation_alias="ESTIMATE_EXTERNAL_SOURCES_BACKGROUND" + ) estimate_geocode_budget_s: float = 12.0 estimate_house_meta_timeout_s: float = 8.0 diff --git a/tradein-mvp/backend/app/services/estimator.py b/tradein-mvp/backend/app/services/estimator.py index 73880ec1..96936eef 100644 --- a/tradein-mvp/backend/app/services/estimator.py +++ b/tradein-mvp/backend/app/services/estimator.py @@ -27,7 +27,7 @@ import math import re import statistics import time -from collections.abc import Callable, Iterable +from collections.abc import Awaitable, Callable, Iterable from dataclasses import dataclass from datetime import UTC, date, datetime, timedelta from typing import Any, Literal @@ -55,6 +55,7 @@ from sqlalchemy import text from sqlalchemy.orm import Session from app.core.config import Settings, settings +from app.core.db import SessionLocal from app.schemas.trade_in import ( AggregatedEstimate, AnalogLot, @@ -862,6 +863,63 @@ def _yandex_valuation_cache_key(address: str, offer_category: str, offer_type: s return hashlib.sha256(payload.encode("utf-8")).hexdigest() +# ── отложенная догрузка внешних источников (фоновый режим) ─────────────────── +# +# Держим ссылки на задачи: asyncio не хранит их сам, и без этого сборщик мусора +# может убить задачу на полпути (документированное поведение create_task). +_DEFERRED_REFRESH_TASKS: set[asyncio.Task[None]] = set() + +# Подобран под текущий прод: 1 воркер uvicorn, mem_limit 768m, max_connections +# 100 у Postgres. Догрузка почти всё время ждёт чужой HTTP, поэтому небольшой +# потолок не занижает пропускную способность прогрева. +_MAX_DEFERRED_REFRESH_TASKS = 8 + + +def _defer_external_refresh(label: str, work: Callable[[Session], Awaitable[object]]) -> None: + """Догрузить внешний источник ПОСЛЕ ответа, наполнив кэш к следующему разу. + + Своя сессия — намеренно: сессия запроса закрывается вместе с ответом, а обе + функции источников делают внутри себя `db.commit()`. Переиспользование чужой + сессии зафиксировало бы её незавершённую работу. + + Все ошибки гасятся: это прогрев кэша, а не часть ответа. Провал означает лишь + то, что следующий запрос по адресу снова промахнётся мимо кэша. + """ + + async def _run() -> None: + db = SessionLocal() + try: + await work(db) + except Exception: + logger.exception("deferred %s: догрузка не удалась (кэш не прогрет)", label) + finally: + db.close() + + # Потолок одновременных догрузок. Без него всплеск по НОВЫМ адресам (ровно + # тот случай, ради которого режим и сделан) породил бы сотни параллельных + # задач, каждая с сессией к базе и HTTP-клиентом — при max_connections=100 + # и mem_limit 768m у backend это отказ вместо ускорения. + # Переполнение не ошибка: прогрев кэша необязателен, следующий запрос по + # адресу попробует снова. + if len(_DEFERRED_REFRESH_TASKS) >= _MAX_DEFERRED_REFRESH_TASKS: + logger.info( + "deferred %s: очередь догрузки заполнена (%d) — пропуск", + label, + _MAX_DEFERRED_REFRESH_TASKS, + ) + return + + try: + task = asyncio.create_task(_run()) + except RuntimeError: + # Нет активного loop (синхронный вызов из теста/скрипта) — молча пропускаем: + # прогрев кэша необязателен по определению. + logger.debug("deferred %s: нет event loop — пропуск", label) + return + _DEFERRED_REFRESH_TASKS.add(task) + task.add_done_callback(_DEFERRED_REFRESH_TASKS.discard) + + async def _get_or_fetch_yandex_valuation_cached( db: Session, *, @@ -869,6 +927,7 @@ async def _get_or_fetch_yandex_valuation_cached( offer_category: str = YANDEX_VALUATION_DEFAULT_CATEGORY, offer_type: str = YANDEX_VALUATION_DEFAULT_TYPE, house_id: int | None = None, + fetch_on_miss: bool = True, ) -> YandexValuationResult | None: """Cached Yandex Valuation lookup. TTL 24h via external_valuations table. @@ -922,6 +981,12 @@ async def _get_or_fetch_yandex_valuation_cached( except Exception as e: logger.warning("yandex_valuation: cache deserialize failed — refetching: %s", e) + # fetch_on_miss=False — режим «только кэш» (фоновый режим внешних источников): вызывающий код + # не хочет ждать внешний HTTP в запросе и сам поставит догрузку в фон. + if not fetch_on_miss: + logger.info("yandex_valuation: cache MISS key=%s — fetch отложен", cache_key[:8]) + return None + # Fresh fetch try: async with YandexValuationScraper( @@ -4285,13 +4350,24 @@ async def estimate_quality( # ── Stage 8: Yandex Valuation as on-demand source (anonymous, cached 24h) ── yandex_val: YandexValuationResult | None = None if geo is not None and geo.full_address: + _y_addr, _y_house = geo.full_address, target_house_id + _bg = settings.estimate_external_sources_background yandex_val = await _with_budget( _get_or_fetch_yandex_valuation_cached( - db, address=geo.full_address, house_id=target_house_id + db, address=_y_addr, house_id=_y_house, fetch_on_miss=not _bg ), settings.estimate_yandex_valuation_timeout_s, label="yandex_valuation", ) + if _bg and yandex_val is None: + # Промах кэша: ответ не ждёт ~6 с внешнего HTTP, догрузка уходит в фон + # и наполняет кэш к следующему обращению по этому адресу. + _defer_external_refresh( + "yandex_valuation", + lambda bg_db: _get_or_fetch_yandex_valuation_cached( + bg_db, address=_y_addr, house_id=_y_house + ), + ) if yandex_val is not None: saved_hist = await asyncio.to_thread(_save_yandex_history_items, db, yandex_val) logger.info( @@ -4311,24 +4387,30 @@ async def estimate_quality( and payload.floor is not None and payload.total_floors is not None ): + _c_bg = settings.estimate_external_sources_background + _c_kwargs: dict[str, Any] = { + "config": RealScraperConfig(), + "address": geo.full_address, + "total_area": payload.area_m2, + "rooms_count": payload.rooms, + "floor": payload.floor, + "total_floors": payload.total_floors, + "repair_type": "cosmetic", + "deal_type": "sale", + "use_cache": True, + "house_id": target_house_id, + } try: cian_val = await _with_budget( - estimate_via_cian_valuation( - db, - config=RealScraperConfig(), - address=geo.full_address, - total_area=payload.area_m2, - rooms_count=payload.rooms, - floor=payload.floor, - total_floors=payload.total_floors, - repair_type="cosmetic", - deal_type="sale", - use_cache=True, - house_id=target_house_id, - ), + estimate_via_cian_valuation(db, fetch_on_miss=not _c_bg, **_c_kwargs), settings.estimate_cian_valuation_timeout_s, label="cian_valuation", ) + if _c_bg and cian_val is None: + _defer_external_refresh( + "cian_valuation", + lambda bg_db: estimate_via_cian_valuation(bg_db, **_c_kwargs), + ) if cian_val is not None and cian_val.sale_price_rub: logger.info( "cian_valuation: price=%s accuracy=%s house_id=%s", diff --git a/tradein-mvp/backend/tests/test_estimate_external_sources_background.py b/tradein-mvp/backend/tests/test_estimate_external_sources_background.py new file mode 100644 index 00000000..dcf44654 --- /dev/null +++ b/tradein-mvp/backend/tests/test_estimate_external_sources_background.py @@ -0,0 +1,212 @@ +"""Внешние источники оценки не ждутся в запросе, а догружаются в фоне. + +Замер на проде 2026-08-22: расчёт по новому адресу 8-12 с, из них ~6 с ждёт +Yandex и ~1.5 с Cian; собственный расчёт укладывается в секунду. По уже +виденному адресу (кэш 24 ч) — 0.4-0.8 с. + +Ключевое свойство, которое делает переключение безопасным: деградация источника +в None — НЕ новое состояние ответа. Ровно так же ведёт себя таймаут +`estimate_*_valuation_timeout_s`, и этот путь работает в проде сегодня. + +Тесты проверяют три вещи, каждая из которых при поломке молча вернула бы +секунды ожидания в пользовательский путь: + 1. `fetch_on_miss=False` действительно НЕ ходит в сеть при промахе кэша; + 2. попадание в кэш работает одинаково в обоих режимах; + 3. отложенная задача берёт СВОЮ сессию и гасит свои ошибки. +""" + +from __future__ import annotations + +import asyncio +from typing import Any +from unittest.mock import MagicMock + +import pytest + + +class _NoRowsDB: + """Сессия, у которой кэш всегда пуст.""" + + def __init__(self) -> None: + self.commits = 0 + self.closed = False + + def execute(self, *_a: Any, **_kw: Any) -> Any: + result = MagicMock() + result.mappings.return_value.first.return_value = None + return result + + def commit(self) -> None: + self.commits += 1 + + def rollback(self) -> None: + return None + + def close(self) -> None: + self.closed = True + + +# ── Yandex ─────────────────────────────────────────────────────────────────── + + +def test_yandex_cache_only_does_not_touch_network(monkeypatch: pytest.MonkeyPatch) -> None: + """fetch_on_miss=False → None и НИ ОДНОГО обращения к скраперу.""" + from app.services import estimator + + def _boom(*_a: Any, **_kw: Any) -> Any: + raise AssertionError("сеть не должна дёргаться в режиме «только кэш»") + + monkeypatch.setattr(estimator, "YandexValuationScraper", _boom) + + out = asyncio.run( + estimator._get_or_fetch_yandex_valuation_cached( + _NoRowsDB(), # type: ignore[arg-type] + address="Екатеринбург, улица Сурикова, 4", + fetch_on_miss=False, + ) + ) + assert out is None + + +def test_yandex_default_still_fetches(monkeypatch: pytest.MonkeyPatch) -> None: + """Дефолт не изменился: без флага источник по-прежнему грузится в запросе.""" + from app.services import estimator + + called: dict[str, bool] = {} + + class _Scraper: + def __init__(self, *_a: Any, **_kw: Any) -> None: + called["constructed"] = True + + async def __aenter__(self) -> _Scraper: + return self + + async def __aexit__(self, *_a: Any) -> None: + return None + + async def fetch_house_history(self, **_kw: Any) -> None: + return None + + monkeypatch.setattr(estimator, "YandexValuationScraper", _Scraper) + + asyncio.run( + estimator._get_or_fetch_yandex_valuation_cached( + _NoRowsDB(), # type: ignore[arg-type] + address="Екатеринбург, улица Сурикова, 4", + ) + ) + assert called.get("constructed") is True + + +# ── отложенная догрузка ────────────────────────────────────────────────────── + + +def test_deferred_refresh_uses_its_own_session(monkeypatch: pytest.MonkeyPatch) -> None: + """Сессия запроса закрывается вместе с ответом — фон обязан открыть свою. + + Переиспользование чужой сессии зафиксировало бы её незавершённую работу: + обе функции источников делают внутри себя `db.commit()`. + """ + from app.services import estimator + + own = _NoRowsDB() + request_db = _NoRowsDB() + monkeypatch.setattr(estimator, "SessionLocal", lambda: own) + seen: dict[str, Any] = {} + + async def _work(db: Any) -> None: + seen["db"] = db + + async def _drive() -> None: + estimator._defer_external_refresh("test", _work) + await asyncio.sleep(0) + await asyncio.sleep(0) + + asyncio.run(_drive()) + + assert seen["db"] is own, "фоновая задача взяла не свою сессию" + assert seen["db"] is not request_db + assert own.closed is True, "фоновая сессия осталась незакрытой" + + +def test_deferred_refresh_swallows_errors(monkeypatch: pytest.MonkeyPatch) -> None: + """Провал прогрева кэша не должен всплывать: это не часть ответа.""" + from app.services import estimator + + own = _NoRowsDB() + monkeypatch.setattr(estimator, "SessionLocal", lambda: own) + + async def _work(_db: Any) -> None: + raise RuntimeError("внешний сервис лёг") + + async def _drive() -> None: + estimator._defer_external_refresh("test", _work) + await asyncio.sleep(0) + await asyncio.sleep(0) + + asyncio.run(_drive()) # не должно поднять + assert own.closed is True + + +def test_deferred_refresh_without_event_loop_is_noop() -> None: + """Синхронный контекст (скрипт/тест) — прогрев необязателен, не падаем.""" + from app.services import estimator + + async def _work(_db: Any) -> None: + raise AssertionError("не должно вызваться") + + estimator._defer_external_refresh("test", _work) # без запущенного loop + + +def test_task_reference_is_retained(monkeypatch: pytest.MonkeyPatch) -> None: + """Ссылка на задачу удерживается: иначе GC может убить её на полпути.""" + from app.services import estimator + + monkeypatch.setattr(estimator, "SessionLocal", _NoRowsDB) + started = asyncio.Event() + release = asyncio.Event() + + async def _work(_db: Any) -> None: + started.set() + await release.wait() + + async def _drive() -> None: + estimator._defer_external_refresh("test", _work) + await started.wait() + assert len(estimator._DEFERRED_REFRESH_TASKS) == 1 + release.set() + await asyncio.sleep(0) + await asyncio.sleep(0) + assert len(estimator._DEFERRED_REFRESH_TASKS) == 0, "ссылка не убрана после завершения" + + asyncio.run(_drive()) + + +def test_deferred_queue_is_bounded(monkeypatch: pytest.MonkeyPatch) -> None: + """Всплеск по новым адресам не должен породить сотни задач. + + Ровно тот случай, ради которого режим и сделан: у backend один воркер, + mem_limit 768m, у Postgres max_connections 100. Неограниченный create_task + превратил бы ускорение в отказ. + """ + from app.services import estimator + + monkeypatch.setattr(estimator, "SessionLocal", _NoRowsDB) + monkeypatch.setattr(estimator, "_MAX_DEFERRED_REFRESH_TASKS", 3) + release = asyncio.Event() + + async def _work(_db: Any) -> None: + await release.wait() + + async def _drive() -> None: + for _ in range(50): + estimator._defer_external_refresh("test", _work) + await asyncio.sleep(0) + assert len(estimator._DEFERRED_REFRESH_TASKS) <= 3, ( + f"очередь догрузки не ограничена: {len(estimator._DEFERRED_REFRESH_TASKS)} задач" + ) + release.set() + for _ in range(5): + await asyncio.sleep(0) + + asyncio.run(_drive()) diff --git a/tradein-mvp/docker-compose.prod.yml b/tradein-mvp/docker-compose.prod.yml index a2fb8211..a3258f9d 100644 --- a/tradein-mvp/docker-compose.prod.yml +++ b/tradein-mvp/docker-compose.prod.yml @@ -246,6 +246,13 @@ services: - path: ./backend/.env.runtime required: false environment: + # Внешние оценки (Yandex/Cian) не ждать в запросе — догружать в фоне. + # Замер 2026-08-22: расчёт по новому адресу 8-12 с, из них ~6 с Yandex и + # ~1.5 с Cian, собственный расчёт < 1 с. Публикация в РБК 30.08 приведёт + # аудиторию на НОВЫЕ адреса, то есть мимо суточного кэша. + # Деградация источника в None — существующее состояние ответа (так же + # ведёт себя таймаут), контракт API не меняется. + ESTIMATE_EXTERNAL_SOURCES_BACKGROUND: "true" DATABASE_URL: "postgresql+psycopg://${TRADEIN_POSTGRES_USER:-tradein}:${TRADEIN_POSTGRES_PASSWORD}@postgres:5432/tradein" PUBLIC_URL: "https://gendsgn.ru/trade-in" CORS_ORIGINS: '["https://gendsgn.ru"]' diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/valuation.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/valuation.py index 42cf75d7..a611012c 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/valuation.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/valuation.py @@ -143,6 +143,7 @@ async def estimate_via_cian_valuation( repair_type: str = "cosmetic", deal_type: str = "sale", use_cache: bool = True, + fetch_on_miss: bool = True, house_id: int | None = None, listing_id: int | None = None, proxy_provider: ProxyProvider | None = None, @@ -164,6 +165,13 @@ async def estimate_via_cian_valuation( logger.info("Cian valuation cache HIT for cache_key=%s...", cache_key[:12]) return cached + # fetch_on_miss=False — режим «только кэш»: вызывающий код не хочет ждать + # внешний HTTP в запросе и сам поставит догрузку в фон. Возврат None здесь + # идёт по тому же graceful-пути, что таймаут и сетевая ошибка. + if not fetch_on_miss: + logger.info("Cian valuation: cache MISS key=%s... — fetch отложен", cache_key[:12]) + return None + # 2. Load auth cookies из cian_session_cookies cookies = load_session(db, config=config) if cookies is None: