From ac11156f7d1136a9a9d3c307d5c4952626149790 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 5 Sep 2026 22:44:46 +0500 Subject: [PATCH] =?UTF-8?q?fix(scraper-kit):=20=D0=BC=D0=B5=D1=82=D0=BA?= =?UTF-8?q?=D0=B0=20=D0=B4=D1=80=D0=B5=D0=B9=D0=BD=D0=B0=20interrupted=3D1?= =?UTF-8?q?=20=D0=B2=D0=BE=20=D0=B2=D1=81=D0=B5=D1=85=20city-=D1=81=D0=B2?= =?UTF-8?q?=D0=B8=D0=BF=D0=B0=D1=85=20(#3333)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit После #3319 резюм подхватывает 'done'-прогоны только с counters.interrupted (_drained_done в scheduler._resume_decision), но ставил метку ровно один avito_city_sweep. У yandex, cian и newbuilding SIGTERM-drain финализировался чистым 'done' с частичными счётчиками: оборванный деплоем обход неотличим от полного и из резюма выпадал, хотя чекпоинт done_buckets есть у всех трёх (combo-метки / имена якорей / номера страниц). done_buckets в дрейн-payload не добавляю: heartbeat мержит jsonb, уже записанные единицы обхода переживают финализатор, а пустой список у multi-anchor yandex затёр бы унаследованный при claim чекпоинт. --- .../tests/test_3333_drain_mark_all_sweeps.py | 212 ++++++++++++++++++ .../src/scraper_kit/orchestration/pipeline.py | 24 +- 2 files changed, 230 insertions(+), 6 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_3333_drain_mark_all_sweeps.py diff --git a/tradein-mvp/backend/tests/test_3333_drain_mark_all_sweeps.py b/tradein-mvp/backend/tests/test_3333_drain_mark_all_sweeps.py new file mode 100644 index 00000000..80a62dcf --- /dev/null +++ b/tradein-mvp/backend/tests/test_3333_drain_mark_all_sweeps.py @@ -0,0 +1,212 @@ +"""SIGTERM-дрейн помечен `interrupted=1` во ВСЕХ свипах, не только у avito (#3333). + +После #3319 резюм подхватывает 'done'-прогоны только с меткой `interrupted` +(`_drained_done` в scheduler._resume_decision), а ставил её ровно один +avito_city_sweep. У yandex (~2356 стр.), cian (~2960) и newbuilding (~2098) +дрейн финализировался чистым `done` с частичными счётчиками: недоделанный обход +объявлен полным, статус тот же, что у честного, — и из резюма он выпадал. + +Резюм есть у всех трёх (чекпоинт `done_buckets`: combo-метки у yandex, имена +якорей у cian, номера страниц у newbuilding), так что метка не диагностическая: +у каждого есть что подхватывать. Оговорка одна — yandex подхватывает только при +единственном якоре (combo-ключ не содержит якоря); прод-режим ровно такой. + +Анти-цикл (последний тест): метка не должна превратить ЛЮБОЙ 'done' в +резюмируемый — иначе источник больше никогда не обходится целиком. +""" + +from __future__ import annotations + +import os + +# Settings собирается автофикстурой conftest'а и требует database_url — как в +# test_3319_citysweep_checkpoint.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 + +ANCHOR_A = (56.83, 60.60, "ekb-center") +ANCHOR_B = (56.79, 60.63, "ekb-south") + + +class _FakeDb: + """Все UPDATE'ы с counters (heartbeat И финализаторы) складываются по порядку.""" + + def __init__(self) -> None: + self.writes: list[dict[str, Any]] = [] + + def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any: + if params and "counters" in params: + self.writes.append(json.loads(params["counters"])) + return MagicMock() + + def commit(self) -> None: ... + def rollback(self) -> None: ... + + +class _FakeAsyncSession: + def __init__(self, *_a: Any, **_kw: Any) -> None: ... + + async def __aenter__(self) -> _FakeAsyncSession: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + +class _NeverCalledScraper: + """Дрейн срабатывает ДО скрапера: любой запрос тут — сломанный порядок проверок.""" + + def __init__(self, *_a: Any, **_kw: Any) -> None: + self._browser = None + self._cffi = None + self.state_extraction_attempts = 1 + self.state_extraction_failures = 0 + self.request_delay_sec = 0.0 + + 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, + use_proxy_pool_browser=False, + browser_http_endpoint=None, + environment="test", + avito_serp_ok_not_banned=True, + ) + + +@pytest.mark.asyncio +async def test_yandex_drain_is_marked_interrupted() -> None: + """yandex-sweep: дрейн на границе якоря → counters.interrupted == 1.""" + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + enrichment = MagicMock() + enrichment.record_yandex_price_history.return_value = 0 + + with ( + patch.object(pl, "YandexRealtyScraper", _NeverCalledScraper), + patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_yandex_city_sweep( + db, # type: ignore[arg-type] + run_id=3333, + config=_config(), + matcher=MagicMock(), + enrichment=enrichment, + enrich_address=False, + shutdown_requested=lambda: True, + ) + + assert db.writes, "дрейн не оставил ни одной записи counters" + assert db.writes[-1].get("interrupted") == 1, ( + "yandex: оборванный дрейном прогон неотличим от полного обхода — резюм его не возьмёт" + ) + + +@pytest.mark.asyncio +async def test_cian_drain_is_marked_interrupted() -> None: + """cian-sweep: дрейн на границе якоря → counters.interrupted == 1.""" + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + + with ( + patch.object(pl, "CianScraper", _NeverCalledScraper), + patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_cian_city_sweep( + db, # type: ignore[arg-type] + run_id=3333, + config=_config(), + matcher=MagicMock(), + anchors=[ANCHOR_A, ANCHOR_B], + enrich_houses=False, + detail_top_n=0, + request_delay_sec=0.0, + shutdown_requested=lambda: True, + ) + + assert db.writes, "дрейн не оставил ни одной записи counters" + assert db.writes[-1].get("interrupted") == 1, ( + "cian: оборванный дрейном прогон неотличим от полного обхода — резюм его не возьмёт" + ) + + +@pytest.mark.asyncio +async def test_newbuilding_drain_is_marked_interrupted() -> None: + """nb-sweep: дрейн до SERP-фазы → counters.interrupted == 1.""" + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + + with ( + patch.object(pl, "AvitoScraper", _NeverCalledScraper), + patch.object(pl, "AsyncSession", _FakeAsyncSession), + patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_avito_newbuilding_sweep( + db, # type: ignore[arg-type] + run_id=3333, + config=_config(), + matcher=MagicMock(), + pages=4, + request_delay_sec=0.0, + shutdown_requested=lambda: True, + ) + + assert db.writes, "дрейн не оставил ни одной записи counters" + assert db.writes[-1].get("interrupted") == 1, ( + "newbuilding: оборванный дрейном прогон неотличим от полного обхода — резюм его не возьмёт" + ) + + +def _prev_run(counters: dict[str, Any]) -> types.SimpleNamespace: + return types.SimpleNamespace( + prev_id=4707, + prev_status="done", + prev_counters=counters, + same_params=True, + age_h=2.0, + interval_days="7", + ) + + +def test_clean_done_is_not_resumed_but_marked_drain_is() -> None: + """Анти-цикл: подхватывается ТОЛЬКО помеченный дрейн, чистый 'done' — нет. + + Общий на все три провайдера: `_resume_decision` смотрит на counters, а не на + источник, поэтому разбор один. Если бы метка была не нужна для подхвата, все + три правки выше были бы записью в лог ради записи в лог. + """ + from scraper_kit.orchestration.scheduler import _resume_decision + + ckpt = {"done_buckets": ["combo-1"], "resume_chain": 0} + + resume_from, verdict = _resume_decision(_prev_run({**ckpt, "interrupted": 1})) + assert resume_from == 4707, f"помеченный дрейн не подхвачен: {verdict}" + + resume_from, verdict = _resume_decision(_prev_run(ckpt)) + assert resume_from is None, "чистый 'done' — полный обход, подхватывать нечего" + assert verdict["resume_reason"] == "status_done" 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 654927ef..818acf74 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 @@ -2112,8 +2112,14 @@ async def run_avito_newbuilding_sweep( logger.info( "nb-sweep run_id=%d: SIGTERM-drain — stopping before SERP phase", run_id ) - runs.update_heartbeat(db, run_id, counters.to_dict()) - runs.mark_done(db, run_id, counters.to_dict()) + # #3333: та же метка дрейна, что у avito_city_sweep (_ckpt(interrupted=1)). + # Без неё оборванный деплоем обход неотличим от полного: статус 'done', + # счётчики частичные — и резюм (`_drained_done` в scheduler) его не берёт. + # done_buckets тут не пишем: heartbeat мержит jsonb, уже записанные + # страницы переживают финализатор. + _drain = {**counters.to_dict(), "interrupted": 1} + runs.update_heartbeat(db, run_id, _drain) + runs.mark_done(db, run_id, _drain) return counters # proxy_provider прокинут для консистентности (#2616) — не load-bearing, @@ -2370,8 +2376,11 @@ async def run_yandex_city_sweep( len(_anchors), name, ) - runs.update_heartbeat(db, run_id, counters.to_dict()) - runs.mark_done(db, run_id, counters.to_dict()) + # #3333: та же метка дрейна, что у avito_city_sweep (_ckpt(interrupted=1)) — + # см. там же. done_buckets пишет combo-heartbeat, jsonb-мерж их сохраняет. + _drain = {**counters.to_dict(), "interrupted": 1} + runs.update_heartbeat(db, run_id, _drain) + runs.mark_done(db, run_id, _drain) return counters logger.info( @@ -2985,8 +2994,11 @@ async def run_cian_city_sweep( len(_anchors), name, ) - runs.update_heartbeat(db, run_id, counters.to_dict()) - runs.mark_done(db, run_id, counters.to_dict()) + # #3333: та же метка дрейна, что у avito_city_sweep (_ckpt(interrupted=1)) — + # см. там же. done_buckets пишет end-of-anchor heartbeat, jsonb-мерж хранит. + _drain = {**counters.to_dict(), "interrupted": 1} + runs.update_heartbeat(db, run_id, _drain) + runs.mark_done(db, run_id, _drain) return counters logger.info( -- 2.45.3