From 5b447ec33d753481a47e49961161c8dd35d1f134 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 22 Aug 2026 15:51:31 +0300 Subject: [PATCH] =?UTF-8?q?feat(tradein/estimate):=20=D0=B2=D0=BD=D0=B5?= =?UTF-8?q?=D1=88=D0=BD=D0=B8=D0=B5=20=D0=BE=D1=86=D0=B5=D0=BD=D0=BA=D0=B8?= =?UTF-8?q?=20=D0=BD=D0=B5=20=D0=B6=D0=B4=D1=83=D1=82=D1=81=D1=8F=20=D0=B2?= =?UTF-8?q?=20=D0=B7=D0=B0=D0=BF=D1=80=D0=BE=D1=81=D0=B5=20=E2=80=94=209?= =?UTF-8?q?=20=D1=81=D0=B5=D0=BA=D1=83=D0=BD=D0=B4=20=D0=BF=D1=80=D0=B5?= =?UTF-8?q?=D0=B2=D1=80=D0=B0=D1=89=D0=B0=D1=8E=D1=82=D1=81=D1=8F=20=D0=B2?= =?UTF-8?q?=201?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Замер на проде 2026-08-22, разбивка одного расчёта по логам: 10.659 старт 10.827 дом найден 0.17 с 11.565 аналоги, 49 кандидатов 0.74 с 11.579 ДКП-коридор 0.01 с 17.576 yandex_valuation ← 6.0 с 19.159 cian_valuation ← 1.5 с 19.252 готово итого 8.6 с Семь с половиной секунд из восьми с половиной — ожидание чужих HTTP. Наша база и сам расчёт укладываются в секунду. По уже виденному адресу (кэш 24 ч) — 0.4-0.8 с, по новому — 8-12.4 с. Геокодинг ни при чём: с готовыми координатами те же 7.5-10.8. Публикация в РБК 30.08 приведёт аудиторию на НОВЫЕ адреса, то есть мимо кэша. Масштабирование контейнеров тут не помогает: время уходит на ожидание чужого ответа, а не на наши вычисления. Что сделано: у обоих источников появился режим «только кэш» (fetch_on_miss). При включённом ESTIMATE_EXTERNAL_SOURCES_BACKGROUND запрос делает лишь чтение кэша (локальный запрос на миллисекунды), а свежая загрузка уходит в фон и наполняет кэш к следующему обращению по тому же адресу. Почему это безопасно: деградация источника в None — НЕ новое состояние ответа. Ровно так же ведёт себя таймаут estimate_*_valuation_timeout_s, и этот путь работает в проде сегодня. Контракт API не меняется. Фоновая задача берёт СВОЮ сессию: сессия запроса закрывается вместе с ответом, а обе функции источников делают внутри себя db.commit() — переиспользование чужой сессии зафиксировало бы её незавершённую работу. По той же причине отвергнут наивный asyncio.gather двух источников на одной сессии. Очередь догрузки ограничена восемью задачами. Без потолка всплеск по новым адресам — ровно тот случай, ради которого режим и сделан — породил бы сотни параллельных задач с сессиями и HTTP-клиентами при max_connections 100 и mem_limit 768m у backend, то есть отказ вместо ускорения. Дефолт в коде False: поведение других окружений не меняется. На проде режим включён через docker-compose.prod.yml у сервиса backend. Отдельно НЕ сделано, хотя предлагалось: снижение таймаутов до 4 с. Замер показал, что свежий запрос к Яндексу занимает 6 с — таймаут 4 обрывал бы его почти всегда, кэш бы не наполнялся, и источник оказался бы тихо отключён. Таймаут здесь страховка от патологии, а не регулятор задержки. Тесты: 7 штук на режим «только кэш», собственную сессию, гашение ошибок, удержание ссылки на задачу и потолок очереди. Фальсифицированы — на неизменённом коде падают 5 из 7 (проходит только сторож неизменности дефолта). Смежные тесты оценщика (29 штук: бюджет ЦИАН, клиентские координаты, аудит) зелёные. --- tradein-mvp/backend/app/core/config.py | 20 ++ tradein-mvp/backend/app/services/estimator.py | 112 +++++++-- ...st_estimate_external_sources_background.py | 212 ++++++++++++++++++ tradein-mvp/docker-compose.prod.yml | 7 + .../scraper_kit/providers/cian/valuation.py | 8 + 5 files changed, 344 insertions(+), 15 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_estimate_external_sources_background.py 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: -- 2.45.3