Some checks failed
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Failing after 8m10s
CI Trade-In / changes (pull_request) Successful in 18s
CI / changes (pull_request) Successful in 34s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
С 13.09 21:14 по 16.09 18:52 пять объявлений (10776456, 10775964, 10775945, 10775747, 10775594), стоящих в очереди подряд, в каждом прогоне приходили пустой SPA-оболочкой ~381 КБ без блока контактов. Счётчика попыток на объявлении не было, они возвращались в голову каждого снапшота и пятью недогрузами подряд выбивали брейкер: 16 прогонов закончились ровно attempted=5 incomplete=5. Миграция 321: listings.detail_incomplete_count / detail_incomplete_at. yandex_detail_backfill пишет недогруз на объявление; снапшот не берёт карточку сутки и не берёт совсем после трёх недогрузов (счётчик incomplete_given_up); повторный недогруз уже известной карточки не двигает брейкер, первый — двигает, защита от системного недогруза сохранена. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
230 lines
10 KiB
Python
230 lines
10 KiB
Python
"""Вечно недогружаемые карточки Яндекса не держат очередь добора (#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 = "<html>" + "x" * 381_157 + "</html>"
|
||
FULL_HTML = "<html>" + "x" * 3_920_119 + '"encryptedPhones":["a"]</html>'
|
||
|
||
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()
|