diff --git a/tradein-mvp/backend/app/tasks/yandex_newbuilding_sweep.py b/tradein-mvp/backend/app/tasks/yandex_newbuilding_sweep.py index ba1c8d73..f54537c8 100644 --- a/tradein-mvp/backend/app/tasks/yandex_newbuilding_sweep.py +++ b/tradein-mvp/backend/app/tasks/yandex_newbuilding_sweep.py @@ -129,10 +129,30 @@ _SELECT_PENDING_HOUSES = """ SELECT 1 FROM market.yandex_jk_enrichment e WHERE e.ext_id = hs.ext_id ) ) - ORDER BY h.id, hs.ext_id NULLS LAST + -- #2924: дом, у которого резолв уже не удался, не берётся повторно раньше чем + -- через :retry_days (образец — geocode_tried_at, задача geocode_missing). + AND ( + CAST(:force AS boolean) = TRUE + OR h.yandex_jk_resolve_tried_at IS NULL + OR h.yandex_jk_resolve_tried_at + < NOW() - make_interval(days => CAST(:retry_days AS integer)) + ) + -- #2924: неудавшиеся — в КОНЕЦ очереди, а не в начало: иначе пять первых по h.id + -- занимали все пять слотов каждую неделю. + ORDER BY h.id, h.yandex_jk_resolve_tried_at NULLS FIRST, hs.ext_id NULLS LAST LIMIT :lim """ +# #2924: сколько дней не повторять резолв дома, у которого он не удался. Семь — как у +# geocode_tried_at: разметка Яндекса могла поменяться, дом мог появиться в SERP позже. +RESOLVE_RETRY_DAYS = 7 + +_MARK_RESOLVE_TRIED = """ + UPDATE houses + SET yandex_jk_resolve_tried_at = NOW() + WHERE id = CAST(:hid AS bigint) +""" + _UPDATE_SLUG = """ UPDATE houses SET yandex_jk_slug = CAST(:slug AS text) @@ -251,7 +271,14 @@ async def enrich_yandex_newbuilding_sweep( result.fetchable = sizing["fetchable"] result.pending = sizing["pending"] - rows = db.execute(text(_SELECT_PENDING_HOUSES), {"force": force, "lim": limit}).mappings().all() + rows = ( + db.execute( + text(_SELECT_PENDING_HOUSES), + {"force": force, "lim": limit, "retry_days": RESOLVE_RETRY_DAYS}, + ) + .mappings() + .all() + ) logger.info( "yandex-nb-sweep: total=%d fetchable=%d pending=%d; " @@ -321,6 +348,7 @@ async def enrich_yandex_newbuilding_sweep( if not resolved: logger.warning("slug unresolved house_id=%s ext_id=%s — skip", house_id, ext_id) result.failed_resolve += 1 + _mark_resolve_tried(db, house_id) continue # Persist slug под SAVEPOINT @@ -335,6 +363,7 @@ async def enrich_yandex_newbuilding_sweep( "persist slug failed house_id=%s ext_id=%s: %s", house_id, ext_id, exc ) result.failed_resolve += 1 + _mark_resolve_tried(db, house_id) continue jk_slug = resolved @@ -448,6 +477,22 @@ async def enrich_yandex_newbuilding_sweep( return result +def _mark_resolve_tried(db: Session, house_id: int) -> None: + """#2924: зафиксировать неудачную попытку резолва — иначе дом снова первый в очереди. + + Отдельный SAVEPOINT: отметка не должна ронять прогон, а её потеря — не фатальна + (дом просто вернётся в очередь на следующем прогоне, как и раньше). + """ + sp = db.begin_nested() + try: + db.execute(text(_MARK_RESOLVE_TRIED), {"hid": house_id}) + sp.commit() + db.commit() + except Exception as exc: + sp.rollback() + logger.warning("mark resolve tried failed house_id=%s: %s", house_id, exc) + + async def _sleep_with_jitter(delay: float, idx: int, total: int, *, force: bool = False) -> None: """Polite anti-bot sleep с ±20% jitter. diff --git a/tradein-mvp/backend/data/sql/269_houses_yandex_jk_resolve_tried_at.sql b/tradein-mvp/backend/data/sql/269_houses_yandex_jk_resolve_tried_at.sql new file mode 100644 index 00000000..4d696b4f --- /dev/null +++ b/tradein-mvp/backend/data/sql/269_houses_yandex_jk_resolve_tried_at.sql @@ -0,0 +1,13 @@ +-- #2924: неудача резолва yandex_jk_slug ничего не помечала — выборка yandex_newbuilding_sweep +-- каждый прогон брала те же дома в том же порядке (ORDER BY h.id), и пять первых +-- занимали все пять слотов (limit 5) неделю за неделей. Замер 19.08: 397 ожидающих, +-- первые пять по h.id — 6706, 6741, 6757, 6774, 6785 — одни и те же в каждом прогоне. +-- +-- Маркер «пробовали резолвить» — по образцу listings.geocode_tried_at (миграция 005): +-- выборка повторяет попытку не раньше чем через N дней и ставит неудавшиеся в конец +-- очереди (ORDER BY tried_at NULLS FIRST), а не в начало. +-- +-- Идемпотентно. Колонка nullable, без дефолта — существующие строки = «не пробовали». +ALTER TABLE houses ADD COLUMN IF NOT EXISTS yandex_jk_resolve_tried_at timestamptz; +COMMENT ON COLUMN houses.yandex_jk_resolve_tried_at IS + '#2924: когда yandex_newbuilding_sweep последний раз пытался разрезолвить yandex_jk_slug и не смог; NULL = не пробовали или резолв удался'; diff --git a/tradein-mvp/backend/tests/test_2924_yandex_resolve_tried_at.py b/tradein-mvp/backend/tests/test_2924_yandex_resolve_tried_at.py new file mode 100644 index 00000000..b196303e --- /dev/null +++ b/tradein-mvp/backend/tests/test_2924_yandex_resolve_tried_at.py @@ -0,0 +1,180 @@ +"""Неудача резолва yandex_jk_slug помечает дом, а очередь не упирается в те же пять (#2924). + +`_SELECT_PENDING_HOUSES` брал ожидающие дома `ORDER BY h.id` и уходил из очереди только +попавший в витрину. Неудачная попытка не писала ничего — ни времени, ни счётчика, — и +следующий прогон брал те же дома в том же порядке. Замер 19.08: 397 ожидающих, первые +пять по h.id (6706, 6741, 6757, 6774, 6785) занимали все пять слотов каждый прогон. + +Маркер — `houses.yandex_jk_resolve_tried_at` (миграция 269), по образцу +`listings.geocode_tried_at`: выборка не берёт дом раньше чем через RESOLVE_RETRY_DAYS, +а неудавшиеся идут в КОНЕЦ очереди (`NULLS FIRST`). + +Тесты идут через реальный `enrich_yandex_newbuilding_sweep` с двойником сессии: он отдаёт +заданные строки на SELECT и запоминает весь выполненный SQL. Резолвер патчится на None — +это «неудача», ради которой маркер и нужен. На origin/main маркер не пишется ни разу. +""" + +from __future__ import annotations + +import os + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") + +from contextlib import contextmanager +from typing import Any +from unittest.mock import AsyncMock, MagicMock, patch + +import pytest + + +class _Result: + def __init__(self, rows: list[dict], scalar: Any = 0) -> None: + self._rows = rows + self._scalar = scalar + + def mappings(self) -> _Result: + return self + + def all(self) -> list[dict]: + return self._rows + + def fetchone(self) -> Any: + return self._rows[0] if self._rows else None + + def first(self) -> Any: + return self._rows[0] if self._rows else None + + def scalar(self) -> Any: + return self._scalar + + def scalar_one(self) -> Any: + return self._scalar + + def scalar_one_or_none(self) -> Any: + return self._scalar + + +class _Session: + """Отдаёт pending-строки на SELECT и запоминает все (sql, params).""" + + def __init__(self, pending: list[dict]) -> None: + self._pending = pending + self.sql: list[tuple[str, dict]] = [] + + def execute(self, statement: Any = None, params: Any = None, *a: Any, **kw: Any) -> _Result: + text_ = str(statement) + self.sql.append((text_, dict(params or {}))) + # Выборка ожидающих — единственный запрос с DISTINCT ON (h.id); COUNT-запросы + # тоже читают houses h + yandex_realty_nb, их путать с выборкой нельзя. + if "DISTINCT ON (h.id)" in text_: + return _Result(self._pending) + return _Result([]) + + def begin_nested(self) -> Any: + sp = MagicMock() + + @contextmanager + def _cm(): + yield sp + + # код зовёт и `sp = db.begin_nested(); sp.commit()` и `with db.begin_nested():` + sp.__enter__ = lambda *_: sp + sp.__exit__ = lambda *_: False + return sp + + def commit(self) -> None: + pass + + def rollback(self) -> None: + pass + + +async def _run(pending: list[dict], *, force: bool = False, resolved: str | None = None): + from app.tasks import yandex_newbuilding_sweep as mod + + db = _Session(pending) + # Резолвер и конфиг импортируются ВНУТРИ функции (late import) — патчим модуль-источник. + with ( + patch( + "scraper_kit.providers.yandex.newbuilding.resolve_yandex_jk_slug", + AsyncMock(return_value=resolved), + ), + patch.object(mod, "_sleep_with_jitter", AsyncMock(return_value=None)), + patch("app.services.scraper_adapters.RealScraperConfig", MagicMock()), + ): + await mod.enrich_yandex_newbuilding_sweep( + db, city="ekaterinburg", limit=5, force=force, request_delay_sec=0 + ) + return db + + +def _tried_updates(db: _Session) -> list[int]: + return [ + int(p["hid"]) + for s, p in db.sql + if "UPDATE houses" in s and "yandex_jk_resolve_tried_at" in s and "hid" in p + ] + + +def _select_sql(db: _Session) -> str: + for s, _ in db.sql: + if "DISTINCT ON (h.id)" in s: + return s + raise AssertionError("выборка ожидающих домов не выполнялась") + + +def _select_params(db: _Session) -> dict: + for s, p in db.sql: + if "DISTINCT ON (h.id)" in s: + return p + raise AssertionError("выборка ожидающих домов не выполнялась") + + +@pytest.mark.asyncio +async def test_failed_resolve_marks_the_house() -> None: + """Головной: резолв не удался → дом получает yandex_jk_resolve_tried_at. + + На origin/main такого UPDATE нет вовсе — дом останется первым в очереди навсегда. + """ + db = await _run( + [{"house_id": 6706, "yandex_jk_slug": None, "yandex_jk_id": None, "ext_id": "286394"}] + ) + assert _tried_updates(db) == [ + 6706 + ], f"неудача резолва не помечена; выполненный SQL: {[s[:50] for s, _ in db.sql]}" + + +@pytest.mark.asyncio +async def test_successful_resolve_does_not_mark() -> None: + """Контроль: удачный резолв маркер НЕ ставит — иначе удачные дома тоже уходили + бы в конец очереди и ждали бы неделю до обогащения.""" + db = await _run( + [{"house_id": 6706, "yandex_jk_slug": None, "yandex_jk_id": None, "ext_id": "286394"}], + resolved="uspenskij", + ) + assert _tried_updates(db) == [], "маркер поставлен при удачном резолве" + + +@pytest.mark.asyncio +async def test_selection_skips_recently_tried_and_pushes_tried_to_the_end() -> None: + """Выборка: не берёт дом, у которого попытка свежее RESOLVE_RETRY_DAYS, и ставит + пробованные в конец (NULLS FIRST). Проверяется по тексту запроса, который реально + ушёл в БД, и по параметрам — retry_days доезжает, а не захардкожен в строке.""" + db = await _run([]) + sql = _select_sql(db) + assert "yandex_jk_resolve_tried_at IS NULL" in sql, "нет фильтра по маркеру" + assert ":retry_days" in sql, "срок повтора не параметризован" + assert "yandex_jk_resolve_tried_at NULLS FIRST" in sql, "пробованные не уходят в конец очереди" + assert _select_params(db).get("retry_days") == 7, _select_params(db) + + +@pytest.mark.asyncio +async def test_force_bypasses_the_marker() -> None: + """Контроль: force=True игнорирует маркер — ручной полный проход остаётся возможным.""" + db = await _run([], force=True) + sql = _select_sql(db) + # Тот же OR-блок, что и у гейта «уже обогащён»: CAST(:force AS boolean) = TRUE OR … + assert ( + sql.count("CAST(:force AS boolean) = TRUE") >= 2 + ), "маркер не обходится через force — второго OR-блока с :force нет" + assert _select_params(db).get("force") is True