From c7a03495e2b5ae42bfe03f176121424e0052b1bf Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 17 Sep 2026 13:54:55 +0500 Subject: [PATCH] =?UTF-8?q?=D0=A1=D0=B2=D0=B8=D0=BF=20=D0=94=D0=BE=D0=BC?= =?UTF-8?q?=D0=9A=D0=BB=D0=B8=D0=BA=D0=B0:=20=D0=BA=D0=BE=D1=80=D0=B7?= =?UTF-8?q?=D0=B8=D0=BD=D0=B0=20=D0=B1=D0=B5=D0=B7=20=D1=81=D0=BE=D1=85?= =?UTF-8?q?=D1=80=D0=B0=D0=BD=D1=91=D0=BD=D0=BD=D1=8B=D1=85=20=D1=81=D1=82?= =?UTF-8?q?=D1=80=D0=BE=D0=BA=20=D0=BD=D0=B5=20=D0=BF=D0=BE=D0=BF=D0=B0?= =?UTF-8?q?=D0=B4=D0=B0=D0=B5=D1=82=20=D0=B2=20=D1=87=D0=B5=D0=BA=D0=BF?= =?UTF-8?q?=D0=BE=D0=B8=D0=BD=D1=82=20(#2406)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Ревью PR #3561: тест таймаута фазы закреплял как контракт done_buckets == ['st']. Лоты ДомКлика копятся в памяти и пишутся одним save_listings после всех корзин, поэтому снятая watchdog'ом (или упавшая на save) фаза не сохраняет ничего, а чекпоинт всё равно получал completed_buckets живого скрейпера — следующий прогон пропускал корзину навсегда (механизм миграции 308). Чекпоинт пополняется только если фаза дошла до конца (флаг _saved после save). Тест развёрнут: при таймауте fetch_city и при падении save_listings done_buckets == [], статус failed. Контроль «полный проход несёт унаследованное ∪ пройденное» — test_3369. Co-Authored-By: Claude Opus 5 --- .../tests/test_2406_sweep_edge_cases.py | 32 ++++++++++++++++--- .../src/scraper_kit/orchestration/pipeline.py | 13 ++++++-- 2 files changed, 39 insertions(+), 6 deletions(-) diff --git a/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py b/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py index 8b15a43d..feeac535 100644 --- a/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py +++ b/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py @@ -3,7 +3,8 @@ Уже покрыто: avito — таймаут якоря и дрейн после первого якоря (test_3319); дрейн ДО первого якоря у yandex/cian/newbuilding (test_3333); отмена у domclick (test_3369). Здесь остаток: - (а) таймаут якоря у cian и yandex, таймаут фазы у domclick; + (а) таймаут якоря у cian и yandex; у domclick снятая/упавшая фаза не отмечает + корзины, строки которых не сохранены; (б) SIGTERM-дрейн ПОСЛЕ первого якоря у cian и yandex; (в) пользовательская отмена (runs.is_cancelled=True) после первого якоря у avito, cian и yandex. @@ -168,7 +169,16 @@ async def test_anchor_timeout_skips_anchor_and_continues(kind: str) -> None: assert db.counters["errors_count"] >= 1 -async def test_domclick_phase_timeout_keeps_completed_buckets_and_is_not_done() -> None: +@pytest.mark.parametrize("broken", ["fetch_timeout", "save_failed"]) +async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> None: + """Корзина без сохранённых строк не попадает в чекпоинт. + + Лоты Домклика копятся в памяти и пишутся ОДНИМ save_listings после всех корзин. + Фаза, снятая watchdog'ом (или упавшая на save), не сохранила ничего — даже из + корзины 'st', которую скрейпер успел пройти. Отметить её пройденной значило бы, + что следующий прогон пропустит её через skip_buckets навсегда (миграция 308). + """ + class _Dc(_Scraper): blocked = False geo_filtered = fetch_errors = bucket_start_index = 0 @@ -176,11 +186,22 @@ async def test_domclick_phase_timeout_keeps_completed_buckets_and_is_not_done() completed_buckets = ["st"] # noqa: RUF012 async def fetch_city(self, **_kw: Any) -> list[Any]: - raise TimeoutError + if broken == "fetch_timeout": + raise TimeoutError + return [object()] + + saved: list[int] = [] + + def _save(_db: Any, lots: list[Any], **_kw: Any) -> tuple[int, int]: + if broken == "save_failed": + raise RuntimeError("save_listings упал") + saved.append(len(lots)) + return len(lots), 0 db = _FakeDb() with ( patch.object(pl, "DomClickScraper", _Dc), + patch.object(pl, "save_listings", _save), patch.object(pl.runs, "is_cancelled", lambda *_a: False), ): counters = await pl.run_domclick_city_sweep( @@ -192,8 +213,11 @@ async def test_domclick_phase_timeout_keeps_completed_buckets_and_is_not_done() request_delay_sec=0.0, ) + assert saved == [], "тест не о том: строки сохранились" assert counters.errors_count >= 1 - assert db.counters["done_buckets"] == ["st"] + assert db.counters["done_buckets"] == [], ( + f"корзины без единой сохранённой строки в чекпоинте: {db.counters['done_buckets']}" + ) assert db.status == "failed" diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py index 4fae93a3..0c69733b 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py @@ -4992,10 +4992,13 @@ async def run_domclick_city_sweep( ) lots: list[ScrapedLot] = [] + # #2406: собранное сохраняется ОДНИМ save_listings после всех корзин, поэтому + # корзина считается пройденной только если фаза дошла до конца (см. чекпоинт ниже). + _saved = False async def _domclick_phase() -> None: """Единственная citywide-фаза: fetch_city + save.""" - nonlocal lots + nonlocal lots, _saved async with DomClickScraper( config, proxy_provider=proxy_provider, @@ -5041,6 +5044,7 @@ async def run_domclick_city_sweep( ) counters.lots_inserted += inserted counters.lots_updated += updated + _saved = True try: await asyncio.wait_for(_domclick_phase(), timeout=_sweep_timeout) @@ -5094,7 +5098,12 @@ async def run_domclick_city_sweep( # #3118: чекпоинт = унаследованное ∪ завершённое в этом прогоне. # Пишем heartbeat'ом СЕЙЧАС (мерж jsonb) — финализаторы ключ не # затирают, и оборванный болезнью финализации прогон его не теряет. - _checkpoint = sorted(skip_buckets | set(_s.completed_buckets)) + # #2406: только если лоты сохранены. Снятая watchdog'ом или упавшая фаза + # не дошла до save_listings — у её корзин в БД ноль строк, а отметка + # «пройдена» заставила бы следующий прогон пропустить их навсегда + # (механизм разобран в миграции 308, из-за него выключены свипы 77/50). + if _saved: + _checkpoint = sorted(skip_buckets | set(_s.completed_buckets)) runs.update_heartbeat(db, run_id, _payload()) counters.bucket_start_index = _s.bucket_start_index