From f6cf948358fffeaf2da919f13362615de297f481 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 5 Sep 2026 23:59:22 +0500 Subject: [PATCH] =?UTF-8?q?fix(scraper):=20=D0=BF=D0=BE=D0=BC=D0=B5=D1=82?= =?UTF-8?q?=D0=B8=D1=82=D1=8C=20SIGTERM-=D0=B4=D1=80=D0=B5=D0=B9=D0=BD=20i?= =?UTF-8?q?nterrupted=3D1=20=D1=83=20full-load'=D0=BE=D0=B2=20=D0=B8=20dom?= =?UTF-8?q?click-=D1=81=D0=B2=D0=B8=D0=BF=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit После #3330/#3346 резюм берёт 'done' только с counters.interrupted=1, но четыре функции с живым чекпоинтом done_buckets финализировали дрейн чистым 'done': cian/yandex/avito full-load (ветка RuntimeError("shutdown")) и domclick city sweep (дрейн до SERP). Оборванный обход был неотличим от полного, а собранные бакеты никто не подхватывал — у avito это ещё и бан-бюджет (#3315). done_buckets в domclick-payload не добавляем: чекпоинт унаследован при claim (#3074), jsonb-мерж его сохраняет. Closes #3355 --- .../tests/test_3355_drain_mark_full_loads.py | 269 ++++++++++++++++++ .../src/scraper_kit/orchestration/pipeline.py | 38 ++- 2 files changed, 302 insertions(+), 5 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py diff --git a/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py b/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py new file mode 100644 index 00000000..569c7d1b --- /dev/null +++ b/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py @@ -0,0 +1,269 @@ +"""SIGTERM-дрейн помечен `interrupted=1` и у full-load'ов с domclick-свипом (#3355). + +#3330/#3346 закрыли city-свипы, но четыре функции с ЖИВЫМ чекпоинтом `done_buckets` +продолжали финализировать дрейн чистым `done`: cian/yandex/avito full-load (ветка +`RuntimeError("shutdown")`) и domclick city sweep (дрейн до SERP-фазы). Резюм +(`_drained_done` в scheduler._resume_decision) такой прогон не берёт — чекпоинт +пишется в никуда. + +Кто чекпоинт ЧИТАЕТ (метка не диагностическая): + * cian_full_load — `skip_set` из `done_buckets` (ключ "room:lo:hi"), + планировщик зовёт с `resume_run_id=_pick_resume(...)`; + * avito_full_load — то же, ДВА job'а (обычный и exhaustive, #3315); + * domclick_city_sweep— `skip_buckets` из `done_buckets` (ROOM_BUCKETS), тоже + `_pick_resume`; + * yandex_full_load — читает `done_buckets` так же, НО у него нет job'а в + scheduler.py: единственный вход с `resume_run_id` — + админский POST. Метка там честно работает только при + ручном подхвате; автоматического потребителя нет. + +Статус финализации нулевого дрейна фиксируется тестом отдельно: honest-status +гейты (runs.mark_done → _sweep_run_did_nothing/_phase_totally_failed/ +_failed_ratio_too_high) могут увести прогон в 'failed'. Это резюмируемо +('failed' ∈ _RESUME_STATUSES), но статус — факт, а не догадка. + +Анти-цикл («чистый done не резюмируется») уже закрыт в +test_3333_drain_mark_all_sweeps.py — здесь не дублируется. +""" + +from __future__ import annotations + +import os + +# Settings собирается автофикстурой conftest'а и требует database_url — как в +# test_3333_drain_mark_all_sweeps.py, до остальных импортов. +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +import json +import types +from typing import Any +from unittest.mock import MagicMock, patch + +import pytest + + +class _FakeDb: + """Каждый UPDATE с counters запоминается вместе с текстом SQL (для статуса).""" + + def __init__(self) -> None: + self.writes: list[tuple[str, dict[str, Any]]] = [] + + def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: + if params and "counters" in params: + self.writes.append((str(stmt), json.loads(params["counters"]))) + return MagicMock() + + def commit(self) -> None: ... + def rollback(self) -> None: ... + + +def _final(db: _FakeDb) -> tuple[str, dict[str, Any]]: + """Последняя запись counters = финализатор.""" + assert db.writes, "дрейн не оставил ни одной записи counters" + sql, counters = db.writes[-1] + if "status = 'done'" in sql: + return "done", counters + if "status = 'failed'" in sql: + return "failed", counters + return "heartbeat", counters + + +class _DrainAtFirstBucket: + """Любой fetch-метод сразу зовёт on_bucket — настоящий `_on_bucket` и роняет дрейн.""" + + def __init__(self, *_a: Any, **_kw: Any) -> None: + self._browser = None + self._cffi = None + self.request_delay_sec = 0.0 + + async def __aenter__(self) -> _DrainAtFirstBucket: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + def __getattr__(self, name: str) -> Any: + async def _fetch(*_a: Any, **kw: Any) -> Any: + on_bucket = kw.get("on_bucket") + assert on_bucket is not None, f"{name} вызван без on_bucket — тест мимо дрейна" + on_bucket("2:0:5000000", []) + raise AssertionError("_on_bucket не оборвал прогон при shutdown_requested()") + + return _fetch + + +class _NeverCalledScraper: + """Дрейн срабатывает ДО скрапера: любой запрос тут — сломанный порядок проверок.""" + + def __init__(self, *_a: Any, **_kw: Any) -> None: ... + + async def __aenter__(self) -> _NeverCalledScraper: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + def __getattr__(self, name: str) -> Any: + async def _boom(*_a: Any, **_kw: Any) -> Any: + raise AssertionError(f"скрапер вызван при дрейне: {name}") + + return _boom + + +def _config() -> types.SimpleNamespace: + return types.SimpleNamespace( + scraper_fetch_mode="cffi", + scraper_proxy_url=None, + scraper_skip_seen_today=False, + use_proxy_pool_browser=False, + browser_http_endpoint=None, + environment="test", + avito_serp_ok_not_banned=True, + cian_full_load_per_fetch_timeout_s=0.0, + ) + + +@pytest.mark.asyncio +async def test_cian_full_load_drain_is_marked_interrupted() -> None: + """cian full-load: дрейн на границе бакета → counters.interrupted == 1.""" + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + with ( + patch.object(pl, "CianScraper", _DrainAtFirstBucket), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_cian_full_load( + db, # type: ignore[arg-type] + run_id=3355, + config=_config(), + matcher=MagicMock(), + request_delay_sec=0.0, + shutdown_requested=lambda: True, + ) + + _status, counters = _final(db) + assert counters.get("interrupted") == 1, ( + "cian full-load: оборванный дрейном прогон неотличим от полного обхода — " + "резюм его чекпоинт не подхватит" + ) + + +@pytest.mark.asyncio +async def test_yandex_full_load_drain_is_marked_interrupted() -> None: + """yandex full-load: дрейн на границе бакета → counters.interrupted == 1.""" + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + with ( + patch.object(pl, "YandexRealtyScraper", _DrainAtFirstBucket), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_yandex_full_load( + db, # type: ignore[arg-type] + run_id=3355, + config=_config(), + matcher=MagicMock(), + enrichment=MagicMock(), + request_delay_sec=0.0, + shutdown_requested=lambda: True, + ) + + _status, counters = _final(db) + assert counters.get("interrupted") == 1, ( + "yandex full-load: оборванный дрейном прогон неотличим от полного обхода" + ) + + +@pytest.mark.asyncio +async def test_avito_full_load_drain_is_marked_interrupted() -> None: + """avito full-load (exhaustive, #3315): дрейн на границе бакета → interrupted == 1.""" + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + with ( + patch.object(pl, "AvitoScraper", _DrainAtFirstBucket), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_avito_full_load( + db, # type: ignore[arg-type] + run_id=3355, + config=_config(), + matcher=MagicMock(), + request_delay_sec=0.0, + shutdown_requested=lambda: True, + ) + + _status, counters = _final(db) + assert counters.get("interrupted") == 1, ( + "avito full-load: дрейн без метки — бан-бюджет тратится на уже собранные бакеты" + ) + + +@pytest.mark.asyncio +async def test_domclick_sweep_drain_is_marked_interrupted() -> None: + """domclick city sweep: дрейн до SERP-фазы → counters.interrupted == 1.""" + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + with ( + patch.object(pl, "DomClickScraper", _NeverCalledScraper), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_domclick_city_sweep( + db, # type: ignore[arg-type] + run_id=3355, + config=_config(), + matcher=MagicMock(), + request_delay_sec=0.0, + shutdown_requested=lambda: True, + ) + + _status, counters = _final(db) + assert counters.get("interrupted") == 1, ( + "domclick: оборванный дрейном прогон неотличим от полного обхода" + ) + + +@pytest.mark.asyncio +async def test_zero_drain_status_is_resumable() -> None: + """Минор ревью #3346: какой СТАТУС даёт нулевой дрейн — и берёт ли его резюм. + + honest-status гейты в runs.mark_done (#2625/#2700/honest-run-status) вправе + увести прогон в 'failed'. Тест фиксирует ФАКТ, а не ожидание: обе ветки + резюмируемы ('failed' — через _RESUME_STATUSES, 'done' — через метку), но + третьего исхода быть не должно. + """ + from scraper_kit.orchestration import pipeline as pl + from scraper_kit.orchestration.scheduler import _RESUME_STATUSES, _resume_decision + + db = _FakeDb() + with ( + patch.object(pl, "AvitoScraper", _DrainAtFirstBucket), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_avito_full_load( + db, # type: ignore[arg-type] + run_id=3355, + config=_config(), + matcher=MagicMock(), + request_delay_sec=0.0, + shutdown_requested=lambda: True, + ) + + status, counters = _final(db) + assert status == "done", f"нулевой дрейн финализирован как {status!r} (было 'done')" + + prev = types.SimpleNamespace( + prev_id=4707, + prev_status=status, + prev_counters={**counters, "done_buckets": ["2:0:5000000"], "resume_chain": 0}, + same_params=True, + age_h=2.0, + interval_days="7", + ) + resume_from, verdict = _resume_decision(prev) + assert resume_from == 4707, ( + f"дрейн со статусом {status!r} не подхвачен резюмом: {verdict} " + f"(_RESUME_STATUSES={sorted(_RESUME_STATUSES)})" + ) 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 818acf74..9fed62f2 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 @@ -3755,7 +3755,15 @@ async def run_cian_full_load( counters.saved_inserted, counters.saved_updated, ) - runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) + # #3355: та же метка дрейна, что у city-свипов (#3330/#3333). Без неё + # оборванный дрейном обход неотличим от полного: статус тот же 'done', + # а `_drained_done` в scheduler._resume_decision смотрит именно на неё — + # и done_buckets этого прогона не подхватывал никто. + runs.mark_done( + db, + run_id, + {**counters.to_dict(), "done_buckets": sorted(done), "interrupted": 1}, + ) return counters logger.exception("cian-full-load run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) @@ -3997,7 +4005,15 @@ async def run_yandex_full_load( counters.saved_inserted, counters.saved_updated, ) - runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) + # #3355: та же метка дрейна, что у city-свипов (#3330/#3333). Без неё + # оборванный дрейном обход неотличим от полного: статус тот же 'done', + # а `_drained_done` в scheduler._resume_decision смотрит именно на неё — + # и done_buckets этого прогона не подхватывал никто. + runs.mark_done( + db, + run_id, + {**counters.to_dict(), "done_buckets": sorted(done), "interrupted": 1}, + ) return counters logger.exception("yandex-full-load run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) @@ -4263,7 +4279,15 @@ async def run_avito_full_load( counters.saved_inserted, counters.saved_updated, ) - runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) + # #3355: та же метка дрейна, что у city-свипов (#3330/#3333). Без неё + # оборванный дрейном обход неотличим от полного: статус тот же 'done', + # а `_drained_done` в scheduler._resume_decision смотрит именно на неё — + # и done_buckets этого прогона не подхватывал никто. + runs.mark_done( + db, + run_id, + {**counters.to_dict(), "done_buckets": sorted(done), "interrupted": 1}, + ) return counters logger.exception("avito-full-load run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) @@ -4385,8 +4409,12 @@ async def run_domclick_city_sweep( # SIGTERM-drain до SERP-фазы: финализируем run как mark_done(partial), # минуя honest-status (это чистый drain, не QRATOR-блок). logger.info("domclick-sweep run_id=%d: SIGTERM-drain — stopping before SERP", run_id) - runs.update_heartbeat(db, run_id, counters.to_dict()) - runs.mark_done(db, run_id, counters.to_dict()) + # #3355: метка дрейна (см. city-свипы #3330/#3333). done_buckets тут не + # пишем — унаследованный при claim (#3074) чекпоинт лежит в counters, и + # jsonb-мерж heartbeat/финализатора его сохраняет. + _drain = {**counters.to_dict(), "interrupted": 1} + runs.update_heartbeat(db, run_id, _drain) + runs.mark_done(db, run_id, _drain) return counters logger.info(