"""Вечно недогружаемые карточки Яндекса не держат очередь добора (#3191). Прод 13.09 21:14 – 16.09 18:52: пять объявлений (10776456, 10775964, 10775945, 10775747, 10775594), стоящих в очереди подряд, в каждом прогоне приходили пустой SPA-оболочкой ~381 КБ без блока контактов. Пять недогрузов подряд выбивали брейкер consecutive_none=5: 16 прогонов закончились ровно attempted=5 incomplete=5, а «успешные» обрывались, едва дойдя до них. За ними — 20 тыс. необогащённых объявлений. Проверки ПО ЗНАЧЕНИЮ: * мок-часть — какие карточки запрошены и какие счётчики вышли; * живой Postgres (в CI есть, локально skip) — что НАСТОЯЩИЙ снапшот-SELECT после недогруза не берёт карточку сутки, берёт её после паузы и выбрасывает после трёх. """ from __future__ import annotations import os import sys import uuid from typing import Any from unittest.mock import AsyncMock, MagicMock, patch os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") _wp_mock = MagicMock() sys.modules.setdefault("weasyprint", _wp_mock) import pytest # noqa: E402 from sqlalchemy import text # noqa: E402 from app.tasks.yandex_detail_backfill import run_yandex_detail_backfill # noqa: E402 _ASYNC_SESSION = "app.tasks.yandex_detail_backfill.AsyncSession" _PARSE = "app.tasks.yandex_detail_backfill.YandexDetailScraper.parse" _SAVE = "app.tasks.yandex_detail_backfill.save_detail_enrichment" _RUNS = "app.tasks.yandex_detail_backfill.runs_mod" _SLEEP = "app.tasks.yandex_detail_backfill.asyncio.sleep" _RESOLVE_PROXY_URL = "app.tasks.yandex_detail_backfill.resolve_proxy_url" # Оболочка как на проде: 381 КБ, контактов нет. Полная — 3,9 МБ с маркером. SHELL_HTML = "" + "x" * 381_157 + "" FULL_HTML = "" + "x" * 3_920_119 + '"encryptedPhones":["a"]' PERPETUAL = 5 NORMAL = 3 def _resp(html: str) -> MagicMock: resp = MagicMock() resp.status_code = 200 resp.text = html return resp async def _run(db: Any, perpetual_urls: set[str]) -> tuple[Any, list[str]]: """Один прогон; оболочку отдаём по URL, всё остальное — полная карточка.""" requested: list[str] = [] async def get(url: str, **_kw: Any) -> MagicMock: requested.append(url) return _resp(SHELL_HTML if url in perpetual_urls else FULL_HTML) session = AsyncMock() session.get = AsyncMock(side_effect=get) ctx = MagicMock() ctx.__aenter__ = AsyncMock(return_value=session) ctx.__aexit__ = AsyncMock(return_value=None) with ( patch(_ASYNC_SESSION, MagicMock(return_value=ctx)), patch(_PARSE, return_value=MagicMock()), patch(_SAVE, MagicMock(return_value=True)), patch(_RUNS, MagicMock()), patch(_SLEEP, new_callable=AsyncMock), patch(_RESOLVE_PROXY_URL, MagicMock(return_value="http://proxy:3128")), ): counters = await run_yandex_detail_backfill( db, run_id=3191, params={"batch_size": PERPETUAL + NORMAL, "budget_sec": 3600} ) return counters, requested # ── 1. Брейкер: повтор известной карточки не рвёт прогон, первый недогруз — рвёт ── def _mock_db(prior_incomplete: int) -> tuple[MagicMock, set[str], list[str]]: rows = [ { "id": i + 1, "source_url": f"https://realty.yandex.ru/offer/{i + 1}/", "detail_incomplete_count": prior_incomplete if i < PERPETUAL else 0, } for i in range(PERPETUAL + NORMAL) ] db = MagicMock() sel = MagicMock() sel.mappings.return_value.all.return_value = rows sel.one.return_value = MagicMock(url_from_offer_id=0, unenrichable_pending=0) sel.scalar_one.return_value = 0 db.execute.return_value = sel urls = [r["source_url"] for r in rows] return db, set(urls[:PERPETUAL]), urls[PERPETUAL:] @pytest.mark.asyncio async def test_retry_of_known_underloaded_cards_does_not_abort_run() -> None: """Пять карточек, уже недогружавшихся раньше, в голове снапшота — прогон идёт дальше.""" db, perpetual, normal = _mock_db(prior_incomplete=1) counters, requested = await _run(db, perpetual) assert set(normal) <= set(requested), ( f"запрошено {len(requested)} из {PERPETUAL + NORMAL}: повторный недогруз известных " "карточек оборвал прогон раньше, чем очередь дошла до следующих" ) assert (counters.attempted, counters.enriched, counters.incomplete) == (8, 3, 5) @pytest.mark.asyncio async def test_first_underload_series_still_aborts_run() -> None: """Защита от системного недогруза жива: пять ПЕРВЫХ недогрузов подряд — обрыв.""" db, perpetual, normal = _mock_db(prior_incomplete=0) counters, requested = await _run(db, perpetual) assert not set(normal) & set(requested), ( "пять первых недогрузов подряд не оборвали прогон — брейкер системного " "недогруза больше не срабатывает" ) assert (counters.attempted, counters.enriched, counters.incomplete) == (5, 0, 5) # ── 2. Очередь на живом Postgres ───────────────────────────────────────────── def _live_session() -> Any | None: try: from sqlalchemy import create_engine from sqlalchemy.orm import sessionmaker dsn = os.environ.get("TEST_DATABASE_URL") or os.environ.get("DATABASE_URL", "") if not dsn or "localhost:5432/test" in dsn: return None engine = create_engine(dsn, future=True) conn = engine.connect() conn.execute(text("SELECT 1")) conn.close() return sessionmaker(bind=engine, future=True)() except Exception: return None @pytest.mark.skipif(_live_session() is None, reason="no reachable Postgres test DB") @pytest.mark.asyncio async def test_perpetual_underloaded_cards_leave_queue_head() -> None: db = _live_session() tag = f"t3191-{uuid.uuid4().hex[:8]}" try: ids: list[int] = [] urls: list[str] = [] for i in range(PERPETUAL + NORMAL): url = f"https://realty.yandex.ru/offer/3191{uuid.uuid4().int % 10**12}/" # Вечные — самые свежие (голова очереди), как 12.09 15:02 на проде. # Дата в будущем: чужие необогащённые строки тестовой БД не встанут впереди. row_id = db.execute( text( """ INSERT INTO listings (source, source_url, source_id, dedup_hash, price_rub, is_active, scraped_at) VALUES ('yandex', :url, :sid, :sid, 5000000, true, TIMESTAMPTZ '2100-01-01' - make_interval(mins => :i)) RETURNING id """ ), {"url": url, "sid": f"{tag}-{i}", "i": i}, ).scalar_one() ids.append(row_id) urls.append(url) db.commit() perpetual = set(urls[:PERPETUAL]) normal = urls[PERPETUAL:] # Прогон 1: карточки ещё не недогружались — первый недогруз, обрыв на пятой. _c1, req1 = await _run(db, perpetual) assert set(req1) == perpetual, req1 # Прогон 2 сразу следом: голова очереди больше не держит — берутся следующие. c2, req2 = await _run(db, perpetual) assert not set(req2) & perpetual, ( "недогруженные час назад карточки снова в снапшоте — прогон опять " "упрётся в них, как 16 прогонов 13.09–16.09" ) assert set(normal) <= set(req2) assert c2.enriched == NORMAL # Сутки спустя: повтор (разовый недогруз мог пройти) и прогон не обрывается. def age(hours: int) -> None: db.execute( text( "UPDATE listings SET detail_incomplete_at = now() - make_interval(hours => :h) " "WHERE id = ANY(:ids)" ), {"h": hours, "ids": ids[:PERPETUAL]}, ) db.commit() age(25) c3, req3 = await _run(db, perpetual) assert perpetual <= set(req3), "после суток паузы карточки обязаны вернуться в очередь" assert set(normal) <= set(req3), "повтор известных карточек оборвал прогон" assert c3.incomplete == PERPETUAL age(25) await _run(db, perpetual) # третья попытка age(25) c5, req5 = await _run(db, perpetual) assert not set(req5) & perpetual, "после трёх недогрузов карточка всё ещё в очереди" assert c5.incomplete_given_up >= PERPETUAL, ( f"incomplete_given_up={c5.incomplete_given_up}: выбывшие из очереди не посчитаны" ) counts = db.execute( text("SELECT detail_incomplete_count FROM listings WHERE id = ANY(:ids) ORDER BY id"), {"ids": ids}, ).scalars() assert list(counts) == [3] * PERPETUAL + [0] * NORMAL finally: db.rollback() db.execute(text("DELETE FROM listings WHERE source_id LIKE :p"), {"p": f"{tag}-%"}) db.commit() db.close()