From b796702d9a62fc002bf2f23648e2dac0f5a2d15f Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sun, 6 Sep 2026 00:51:47 +0500 Subject: [PATCH 1/2] =?UTF-8?q?fix(domclick):=20=D1=87=D0=B5=D0=BA=D0=BF?= =?UTF-8?q?=D0=BE=D0=B8=D0=BD=D1=82=20done=5Fbuckets=20=D0=B2=D0=BE=20?= =?UTF-8?q?=D0=B2=D1=81=D0=B5=D1=85=20payload'=D0=B0=D1=85=20=D1=81=D0=B2?= =?UTF-8?q?=D0=B8=D0=BF=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit cancel-ветка run_domclick_city_sweep писала голый counters.to_dict(); писатель — app-level scrape_runs с полной перезаписью (CAST(:counters AS jsonb)), а 'cancelled' входит в _RESUME_STATUSES → отменённый прогон закрывался с пустым чекпоинтом и резюм пересобирал все корзины. Те же потери были у финального heartbeat и у mark_banned/mark_failed/mark_done (banned/failed тоже резюмируемы): они затирали чекпоинт, записанный после SERP-фазы. Один _payload() на функцию — done_buckets едет в каждой записи counters. Refs #3355, PR #3363. Closes #3369 --- .../test_3369_domclick_cancel_checkpoint.py | 133 ++++++++++++++++++ .../src/scraper_kit/orchestration/pipeline.py | 32 +++-- 2 files changed, 156 insertions(+), 9 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_3369_domclick_cancel_checkpoint.py diff --git a/tradein-mvp/backend/tests/test_3369_domclick_cancel_checkpoint.py b/tradein-mvp/backend/tests/test_3369_domclick_cancel_checkpoint.py new file mode 100644 index 00000000..98f47e05 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3369_domclick_cancel_checkpoint.py @@ -0,0 +1,133 @@ +"""Cancel-ветка domclick-свипа обязана нести чекпоинт `done_buckets` (#3369). + +Доводка #3355/PR #3363: дрейн- и бан-ветки `run_domclick_city_sweep` кладут +`done_buckets` в свой payload явно, а cancel-ветка писала голый +`counters.to_dict()`. Мержа jsonb тут нет: писатель — app-level +`app/services/scrape_runs.py` (`counters = CAST(:counters AS jsonb)`, ПОЛНАЯ +перезапись), а scheduler при claim `done_buckets` не наследует. `cancelled` +входит в `_RESUME_STATUSES` → отменённый прогон закрывался с пустым чекпоинтом, +и резюм пересобирал все шесть корзин заново. + +Проверка по ЗНАЧЕНИЮ (не по факту записи): унаследованные две корзины обязаны +дословно оказаться в payload'е отмены. +""" + +from __future__ import annotations + +import os + +# Settings собирается автофикстурой conftest'а и требует database_url — как в +# test_3355_drain_mark_full_loads.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; SELECT по `rid` отдаёт чекпоинт.""" + + def __init__(self, prev_counters: dict[str, Any] | None = None) -> None: + self.writes: list[tuple[str, dict[str, Any]]] = [] + self._prev_counters = prev_counters + + 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() + if params and "rid" in params: + row = types.SimpleNamespace(counters=self._prev_counters) + return MagicMock(**{"fetchone.return_value": row if self._prev_counters else None}) + return MagicMock() + + def commit(self) -> None: ... + + def rollback(self) -> None: ... + + +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: + raise AssertionError(f"cancel-ветка не сработала: тронут скрапер ({name})") + + +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_domclick_sweep_cancel_keeps_inherited_checkpoint() -> None: + """Отмена до SERP-фазы: payload несёт ровно унаследованные две корзины.""" + from scraper_kit.orchestration import pipeline as pl + + inherited = ["1:0:5000000", "2:0:5000000"] + db = _FakeDb(prev_counters={"done_buckets": inherited}) + with ( + patch.object(pl, "DomClickScraper", _NeverCalledScraper), + patch.object(pl.runs, "is_cancelled", lambda *_a: True), + ): + await pl.run_domclick_city_sweep( + db, # type: ignore[arg-type] + run_id=3369, + config=_config(), + matcher=MagicMock(), + request_delay_sec=0.0, + shutdown_requested=lambda: False, + resume_run_id=3368, + ) + + assert db.writes, "cancel не оставил ни одной записи counters" + _sql, counters = db.writes[-1] + assert counters.get("done_buckets") == inherited, ( + f"cancel: чекпоинт потерян ({counters.get('done_buckets')!r} вместо " + f"{inherited!r}) — резюм пересоберёт уже собранные корзины" + ) + + +@pytest.mark.asyncio +async def test_domclick_sweep_cancel_marker_is_not_resume_status_only() -> None: + """Отмена без чекпоинта в предшественнике не выдумывает корзин (пустой список).""" + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + with ( + patch.object(pl, "DomClickScraper", _NeverCalledScraper), + patch.object(pl.runs, "is_cancelled", lambda *_a: True), + ): + await pl.run_domclick_city_sweep( + db, # type: ignore[arg-type] + run_id=3369, + config=_config(), + matcher=MagicMock(), + request_delay_sec=0.0, + shutdown_requested=lambda: False, + ) + + assert db.writes, "cancel не оставил ни одной записи counters" + _sql, counters = db.writes[-1] + assert counters.get("done_buckets") == [], ( + f"cancel без предшественника: ожидался пустой чекпоинт, получено " + f"{counters.get('done_buckets')!r}" + ) 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 8222e969..0267b433 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 @@ -4426,6 +4426,16 @@ async def run_domclick_city_sweep( # Мутируемый контейнер для захвата scraper-ссылки из замыкания. _scraper_ref: list[DomClickScraper] = [] + # #3369: чекпоинт (done_buckets) обязан ехать в КАЖДОЙ записи counters этой + # функции. Писатель — app-level scrape_runs: `counters = CAST(:counters AS jsonb)`, + # полная перезапись без jsonb-мержа, а scheduler при claim done_buckets не + # наследует. Любая запись без ключа стирает чекпоинт, и резюм (cancelled/banned/ + # failed ∈ _RESUME_STATUSES) пересобирает уже собранные корзины заново. + _checkpoint: list[str] = [] + + def _payload() -> dict[str, Any]: + return {**counters.to_dict(), "done_buckets": _checkpoint} + try: # ── Checkpoint/resume (#3118): корзина ROOM_BUCKETS — единица обхода ── # Вместе со сдвигом #2854 подхват превращает случайную ротацию в @@ -4451,10 +4461,14 @@ async def run_domclick_city_sweep( len(skip_buckets), ) + _checkpoint = sorted(skip_buckets) + # ── Cooperative cancel перед SERP-фазой ────────────────────────────── if runs.is_cancelled(db, run_id): logger.info("domclick-sweep run_id=%d: cancelled before SERP", run_id) - runs.update_heartbeat(db, run_id, counters.to_dict()) + # #3369: 'cancelled' ∈ _RESUME_STATUSES — payload без done_buckets стирал + # унаследованный чекпоинт, и резюм пересобирал все корзины заново. + runs.update_heartbeat(db, run_id, _payload()) return counters elif shutdown_requested(): # SIGTERM-drain до SERP-фазы: финализируем run как mark_done(partial), @@ -4588,13 +4602,13 @@ async def run_domclick_city_sweep( # #3118: чекпоинт = унаследованное ∪ завершённое в этом прогоне. # Пишем heartbeat'ом СЕЙЧАС (мерж jsonb) — финализаторы ключ не # затирают, и оборванный болезнью финализации прогон его не теряет. - _done_now = sorted(skip_buckets | set(_s.completed_buckets)) - runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": _done_now}) + _checkpoint = sorted(skip_buckets | set(_s.completed_buckets)) + runs.update_heartbeat(db, run_id, _payload()) counters.bucket_start_index = _s.bucket_start_index # pages_fetched: worst-case число страниц (buckets × pages cap). counters.pages_fetched = _num_fetches - runs.update_heartbeat(db, run_id, counters.to_dict()) + runs.update_heartbeat(db, run_id, _payload()) # ── ЧЕСТНЫЙ СТАТУС (#1968, ужесточён #2657) ─────────────────────────── # Распознанный QRATOR-блок НИКОГДА не даёт done. Домклик тут структурно @@ -4621,7 +4635,7 @@ async def run_domclick_city_sweep( run_id, f"QRATOR block aborted sweep — {counters.lots_fetched} listings " "collected before abort (#2657)", - counters.to_dict(), + _payload(), # #2687: эта ветка входится ТОЛЬКО при распознанном QRATOR-блоке # (counters.blocked == 1, DomClickBlockedError — площадка показала # block-страницу). Это ровно определение BAN_KIND_PLATFORM @@ -4651,7 +4665,7 @@ async def run_domclick_city_sweep( f"sweep оборван: пройдено {counters.buckets_completed} комнатных " f"бакетов из {counters.buckets_total}, собрано " f"{counters.lots_fetched} лотов (#2670)", - counters.to_dict(), + _payload(), ) elif counters.lots_fetched == 0 and counters.errors_count > 0: logger.error( @@ -4663,10 +4677,10 @@ async def run_domclick_city_sweep( db, run_id, "fetch errors — 0 listings", - counters.to_dict(), + _payload(), ) else: - runs.mark_done(db, run_id, counters.to_dict()) + runs.mark_done(db, run_id, _payload()) logger.info( "domclick-sweep run_id=%d done: lots=%d (ins=%d/upd=%d) " @@ -4686,5 +4700,5 @@ async def run_domclick_city_sweep( except Exception as exc: logger.exception("domclick-sweep run_id=%d: fatal error", run_id) - runs.mark_failed(db, run_id, str(exc), counters.to_dict()) + runs.mark_failed(db, run_id, str(exc), _payload()) raise -- 2.45.3 From 140caa4ddf5c46fcd778498772b0d72ed8aaa664 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sun, 6 Sep 2026 01:13:15 +0500 Subject: [PATCH 2/2] =?UTF-8?q?fix(domclick):=20=D1=87=D0=B5=D1=81=D1=82?= =?UTF-8?q?=D0=BD=D0=B0=D1=8F=20=D0=BC=D0=BE=D1=82=D0=B8=D0=B2=D0=B8=D1=80?= =?UTF-8?q?=D0=BE=D0=B2=D0=BA=D0=B0=20=D1=87=D0=B5=D0=BA=D0=BF=D0=BE=D0=B8?= =?UTF-8?q?=D0=BD=D1=82=D0=B0=20+=20=D0=BE=D0=B4=D0=B8=D0=BD=20=5Fpayload(?= =?UTF-8?q?)=20=D0=BD=D0=B0=20=D0=B2=D1=81=D0=B5=20=D0=B2=D0=B5=D1=82?= =?UTF-8?q?=D0=BA=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Правка ревью к PR #3372: комментарии и текст коммита обещали восстановление потери, которой нет. Факт: pipeline.py импортирует scraper_kit.orchestration.runs, и все четыре его писателя МЕРЖАТ jsonb (`counters = COALESCE(counters,'{}') || CAST(:counters AS jsonb)`, runs.py:708/776/817/880), а _pick_resume наследует done_buckets при claim (scheduler.py:741-748, #3074) — чекпоинт domclick в БД не терялся. Перезаписывающий двойник app/services/scrape_runs.py этой функцией не вызывается. Поэтому done_buckets в каждом payload — единообразие и защита от смены писателя (если запись пойдёт через app-копию или kit-писатель станет перезаписывающим), а не спасение данных. Комментарии переписаны под этот факт. Дрейн-ветка и NoProxyAvailableError строили dict вручную — переведены на {**_payload(), "interrupted": 1} и _payload(), чтобы «один _payload() на функцию» было правдой. Тест: к cancel добавлены ассерты на финальный heartbeat и mark_done — оба payload'а несут унаследованное ∪ пройденное этим прогоном. Refs #3355, PR #3363. Closes #3369 --- .../test_3369_domclick_cancel_checkpoint.py | 79 ++++++++++++++++--- .../src/scraper_kit/orchestration/pipeline.py | 45 +++++------ 2 files changed, 89 insertions(+), 35 deletions(-) diff --git a/tradein-mvp/backend/tests/test_3369_domclick_cancel_checkpoint.py b/tradein-mvp/backend/tests/test_3369_domclick_cancel_checkpoint.py index 98f47e05..9effa69d 100644 --- a/tradein-mvp/backend/tests/test_3369_domclick_cancel_checkpoint.py +++ b/tradein-mvp/backend/tests/test_3369_domclick_cancel_checkpoint.py @@ -1,15 +1,16 @@ -"""Cancel-ветка domclick-свипа обязана нести чекпоинт `done_buckets` (#3369). +"""Каждая запись counters domclick-свипа несёт чекпоинт `done_buckets` (#3369). Доводка #3355/PR #3363: дрейн- и бан-ветки `run_domclick_city_sweep` кладут -`done_buckets` в свой payload явно, а cancel-ветка писала голый -`counters.to_dict()`. Мержа jsonb тут нет: писатель — app-level -`app/services/scrape_runs.py` (`counters = CAST(:counters AS jsonb)`, ПОЛНАЯ -перезапись), а scheduler при claim `done_buckets` не наследует. `cancelled` -входит в `_RESUME_STATUSES` → отменённый прогон закрывался с пустым чекпоинтом, -и резюм пересобирал все шесть корзин заново. +`done_buckets` в payload явно, а cancel-ветка, финальный heartbeat и финализаторы +писали голый `counters.to_dict()`. Это НЕ было потерей данных: writer — +`scraper_kit.orchestration.runs`, все его писатели мержат jsonb +(`counters = COALESCE(counters,'{}') || CAST(:counters AS jsonb)`), а `_pick_resume` +наследует ключ при claim (#3074). Инвариант ставится ради единообразия и +независимости веток от того, какой писатель окажется на другом конце +(перезаписывающий двойник `app/services/scrape_runs.py` существует). -Проверка по ЗНАЧЕНИЮ (не по факту записи): унаследованные две корзины обязаны -дословно оказаться в payload'е отмены. +Проверка по ЗНАЧЕНИЮ, а не по факту записи: чекпоинт обязан дословно лежать в +payload'е отмены, финального heartbeat'а и `mark_done`. """ from __future__ import annotations @@ -101,8 +102,8 @@ async def test_domclick_sweep_cancel_keeps_inherited_checkpoint() -> None: assert db.writes, "cancel не оставил ни одной записи counters" _sql, counters = db.writes[-1] assert counters.get("done_buckets") == inherited, ( - f"cancel: чекпоинт потерян ({counters.get('done_buckets')!r} вместо " - f"{inherited!r}) — резюм пересоберёт уже собранные корзины" + f"cancel: payload без чекпоинта ({counters.get('done_buckets')!r} вместо " + f"{inherited!r}) — ветка зависит от того, мержит ли писатель jsonb" ) @@ -131,3 +132,59 @@ async def test_domclick_sweep_cancel_marker_is_not_resume_status_only() -> None: f"cancel без предшественника: ожидался пустой чекпоинт, получено " f"{counters.get('done_buckets')!r}" ) + + +class _CleanSweepScraper: + """Свип прошёл все корзины без блока: ветка честного статуса → mark_done.""" + + def __init__(self, *_a: Any, **_kw: Any) -> None: + self.blocked = False + self.geo_filtered = 0 + self.fetch_errors = 0 + self.buckets_completed = 1 + self.buckets_total = 1 + self.completed_buckets = ["3:0:5000000"] + self.bucket_start_index = 0 + self.request_delay_sec = 0.0 + + async def __aenter__(self) -> _CleanSweepScraper: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + async def fetch_city(self, **_kw: Any) -> list[Any]: + return [] + + +@pytest.mark.asyncio +async def test_domclick_sweep_final_heartbeat_and_mark_done_carry_checkpoint() -> None: + """Финальный heartbeat и mark_done несут унаследованное ∪ пройденное этим прогоном.""" + from scraper_kit.orchestration import pipeline as pl + + inherited = ["1:0:5000000", "2:0:5000000"] + expected = sorted([*inherited, "3:0:5000000"]) + db = _FakeDb(prev_counters={"done_buckets": inherited}) + with ( + patch.object(pl, "DomClickScraper", _CleanSweepScraper), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_domclick_city_sweep( + db, # type: ignore[arg-type] + run_id=3369, + config=_config(), + matcher=MagicMock(), + request_delay_sec=0.0, + shutdown_requested=lambda: False, + resume_run_id=3368, + ) + + assert len(db.writes) >= 2, f"ожидались heartbeat + финализатор, есть {len(db.writes)}" + sql_done, counters_done = db.writes[-1] + assert "finished_at" in sql_done, f"последняя запись — не финализатор: {sql_done!r}" + _sql_hb, counters_hb = db.writes[-2] + for label, payload in (("финальный heartbeat", counters_hb), ("mark_done", counters_done)): + assert payload.get("done_buckets") == expected, ( + f"{label}: payload без чекпоинта ({payload.get('done_buckets')!r} вместо " + f"{expected!r}) — ветка зависит от того, мержит ли писатель jsonb" + ) 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 0267b433..001ec4a8 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 @@ -4426,11 +4426,16 @@ async def run_domclick_city_sweep( # Мутируемый контейнер для захвата scraper-ссылки из замыкания. _scraper_ref: list[DomClickScraper] = [] - # #3369: чекпоинт (done_buckets) обязан ехать в КАЖДОЙ записи counters этой - # функции. Писатель — app-level scrape_runs: `counters = CAST(:counters AS jsonb)`, - # полная перезапись без jsonb-мержа, а scheduler при claim done_buckets не - # наследует. Любая запись без ключа стирает чекпоинт, и резюм (cancelled/banned/ - # failed ∈ _RESUME_STATUSES) пересобирает уже собранные корзины заново. + # #3369: чекпоинт (done_buckets) едет в КАЖДОЙ записи counters этой функции — + # это единообразие и защита от смены писателя, а НЕ восстановление потери. + # Как есть сегодня: пишет kit-овый scraper_kit.orchestration.runs (импорт выше), + # и все четыре его писателя МЕРЖАТ jsonb — `counters = COALESCE(counters,'{}') + # || CAST(:counters AS jsonb)`, так что ключ, записанный раньше, переживает + # payload без него; плюс _pick_resume наследует done_buckets при claim (#3074). + # То есть чекпоинт в БД не терялся. Перезаписывающий двойник существует — + # app/services/scrape_runs.py (`counters = CAST(:counters AS jsonb)`), — но + # этой функцией не вызывается. Полный payload делает ветки нечувствительными + # к тому, какой из двух писателей окажется на другом конце. _checkpoint: list[str] = [] def _payload() -> dict[str, Any]: @@ -4466,27 +4471,19 @@ async def run_domclick_city_sweep( # ── Cooperative cancel перед SERP-фазой ────────────────────────────── if runs.is_cancelled(db, run_id): logger.info("domclick-sweep run_id=%d: cancelled before SERP", run_id) - # #3369: 'cancelled' ∈ _RESUME_STATUSES — payload без done_buckets стирал - # унаследованный чекпоинт, и резюм пересобирал все корзины заново. + # #3369: 'cancelled' ∈ _RESUME_STATUSES, прогон резюмируем — пишем + # полный payload, как и все прочие ветки (см. _payload выше). runs.update_heartbeat(db, run_id, _payload()) return counters elif shutdown_requested(): # 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) - # #3355: метка дрейна (см. city-свипы #3330/#3333) И явный перенос - # чекпоинта. Никакого jsonb-мержа тут нет: domclick пишет counters - # через app-level scrape_runs (update_heartbeat/mark_done делают - # `counters = CAST(:counters AS jsonb)` — ПОЛНАЯ перезапись), а - # scheduler при claim done_buckets не наследует. Без явного - # done_buckets дрейн-прогон закрылся бы с пустым чекпоинтом и резюм - # пересобирал бы все шесть корзин заново. Тот же приём, что в - # ban-ветке ниже (except NoProxyAvailableError). - _drain = { - **counters.to_dict(), - "done_buckets": sorted(skip_buckets), - "interrupted": 1, - } + # #3355: метка дрейна (см. city-свипы #3330/#3333) поверх обычного + # payload'а. Чекпоинт тут не «спасается»: kit-писатель мержит jsonb + # (см. _payload выше) — done_buckets в payload'е нужен для + # единообразия, а interrupted=1 отличает дрейн от честного done. + _drain = {**_payload(), "interrupted": 1} runs.update_heartbeat(db, run_id, _drain) runs.mark_done(db, run_id, _drain) return counters @@ -4571,16 +4568,16 @@ async def run_domclick_city_sweep( # хотя площадка вообще не была затронута. Тот же диагноз и тот же по духу # обработчик, что у run_avito_full_load/run_cian_full_load/run_yandex_full_load # (см. except NoProxyAvailableError там же) — ban_kind_of_exception() относит тип - # исключения к BAN_KIND_INFRA («наша инфраструктура», не площадка). done_buckets - # сохраняем как унаследованные skip_buckets — этот прогон новых не завершил, но - # терять уже собранный чекпоинт (#2687-класс дефекта) resume не должен. + # исключения к BAN_KIND_INFRA («наша инфраструктура», не площадка). В payload'е + # done_buckets == унаследованные skip_buckets: этот прогон новых корзин не + # завершил (_checkpoint обновляется только после SERP-фазы). logger.error("domclick-sweep run_id=%d: no proxy available — %s", run_id, exc) counters.errors_count += 1 runs.mark_banned( db, run_id, f"domclick sweep aborted: {exc}", - {**counters.to_dict(), "done_buckets": sorted(skip_buckets)}, + _payload(), ban_kind=ban_kind_of_exception(exc), ) return counters -- 2.45.3