From 0ed934ea09984f2cd0fd0adfb92fa8ffc33c6ccd Mon Sep 17 00:00:00 2001 From: bot-backend Date: Wed, 2 Sep 2026 14:44:38 +0500 Subject: [PATCH] =?UTF-8?q?fix(scraper-kit):=20=D1=87=D0=B5=D0=BA=D0=BF?= =?UTF-8?q?=D0=BE=D0=B8=D0=BD=D1=82=20avito=5Fcity=5Fsweep=20=D0=B4=D0=BE?= =?UTF-8?q?=D0=B6=D0=B8=D0=B2=D0=B0=D0=B5=D1=82=20=D0=B4=D0=BE=20=D1=84?= =?UTF-8?q?=D0=B8=D0=BD=D0=B0=D0=BB=D0=B8=D0=B7=D0=B0=D1=82=D0=BE=D1=80?= =?UTF-8?q?=D0=B0=20(#3319)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 0 из 67 прогонов за 60 дней имели done_buckets в counters: точку писала одна строка внутри цикла якорей, а каждый выход (mark_done — включая ранний #1950 «SERP собран, detail заблокирован», — mark_banned, mark_failed) отдавал голый counters.to_dict(). Точка держалась только на jsonb-мерже в runs.py, то есть на свойстве чужого модуля, которого этот файл не проверяет. - payload любого выхода собирается одной функцией _ckpt() — done_buckets несут все 14 записей, а не одна; - якорь, умерший по таймауту, больше не считается пройденным (тот же инвариант, что у generic-except): SERP мог успеть, detail нет, и резюм пропускал такой якорь навсегда при штатно завершившемся прогоне; - SIGTERM-дрейн помечается counters.interrupted=1 и участвует в резюме. Статус остаётся 'done' — ни один читатель статуса не меняется; метка та же, что у rosreestr_dkp-дрейна. 'done' в _RESUME_STATUSES НЕ добавлен: чистый полный обход резюмить нечего. Дрейн перед IMV-фазой помечен отдельно (imv_phase_drained) — якоря там пройдены все, подхват собрал бы ноль. --- .../tests/test_3319_citysweep_checkpoint.py | 214 ++++++++++++++++++ .../src/scraper_kit/orchestration/pipeline.py | 56 +++-- .../scraper_kit/orchestration/scheduler.py | 9 +- 3 files changed, 262 insertions(+), 17 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_3319_citysweep_checkpoint.py diff --git a/tradein-mvp/backend/tests/test_3319_citysweep_checkpoint.py b/tradein-mvp/backend/tests/test_3319_citysweep_checkpoint.py new file mode 100644 index 00000000..6ae05a73 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3319_citysweep_checkpoint.py @@ -0,0 +1,214 @@ +"""Чекпоинт avito_city_sweep доживает до финализатора (#3319). + +Прод-факт, из которого выросла задача: 0 из 67 прогонов за 60 дней имеют в +counters ключ done_buckets. Механизм #3074 (запись точки) и механизм #930 +(подхват точки) существуют оба, но между ними нет ни одного прогона: точку +писала ровно одна строка внутри цикла якорей, а КАЖДЫЙ выход из прогона +(mark_done — включая ранний выход #1950 «SERP собран, detail заблокирован», — +mark_banned, mark_failed) отдавал голый counters.to_dict() без неё. + +Три инварианта, ради которых тест: + 1. done-выход несёт done_buckets — иначе точка существует только в логе. + 2. Якорь, умерший по таймауту, НЕ пройден: SERP мог успеть, detail нет. + Пройденным его записать = резюм пропустит его навсегда и молча. + 3. SIGTERM-дрейн отличим от полного обхода (counters.interrupted=1) и + участвует в резюме — статус у обоих 'done', счётчики частичные. +""" + +from __future__ import annotations + +import os + +# Settings собирается автофикстурой conftest'а и требует database_url. Выставляем +# до остальных импортов — так же, как в test_3074_avito_anchor_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 _FakeScraper: + """Двойник AvitoScraper: помнит визиты, роняет заданный якорь заданной ошибкой.""" + + visited: list[tuple[float, float]] = [] # noqa: RUF012 — тестовый сборник + raise_on: tuple[float, float] | None = None + exc: type[BaseException] | None = None + lots_per_anchor: int = 0 + + def __init__(self, *_a: Any, **_kw: Any) -> None: + self._browser = None + self._cffi = None + + async def fetch_around(self, lat: float, lon: float, *_a: Any, **_kw: Any) -> list: + _FakeScraper.visited.append((lat, lon)) + if _FakeScraper.raise_on == (lat, lon) and _FakeScraper.exc is not None: + raise _FakeScraper.exc("якорь сорвался") + return [MagicMock() for _ in range(_FakeScraper.lots_per_anchor)] + + +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, + ) + + +async def _run( + *, + raise_on: tuple[float, float] | None = None, + exc: type[BaseException] | None = None, + shutdown_after_first: bool = False, + saved: tuple[int, int] = (0, 0), + lots_per_anchor: int = 0, +) -> _FakeDb: + from scraper_kit.orchestration import pipeline as pl + + _FakeScraper.visited = [] + _FakeScraper.raise_on = raise_on + _FakeScraper.exc = exc + _FakeScraper.lots_per_anchor = lots_per_anchor + db = _FakeDb() + + def _shutdown() -> bool: + return shutdown_after_first and bool(_FakeScraper.visited) + + with ( + patch.object(pl, "AvitoScraper", _FakeScraper), + patch.object(pl, "AsyncSession", _FakeAsyncSession), + patch.object(pl, "save_listings", lambda *_a, **_kw: saved), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_avito_city_sweep( + db, # type: ignore[arg-type] + run_id=3319, + config=_config(), + matcher=MagicMock(), + enrichment=MagicMock(), + anchors=[ANCHOR_A, ANCHOR_B], + enrich_houses=False, + enrich_imv=False, + detail_top_n=0, + shutdown_requested=_shutdown, + ) + return db + + +@pytest.mark.asyncio +async def test_done_exit_carries_checkpoint() -> None: + """Финализатор полного обхода несёт done_buckets, а не голые счётчики.""" + db = await _run() + + assert db.writes[-1].get("done_buckets") == ["ekb-center", "ekb-south"], ( + "финальный (done) выход отдал counters без чекпоинта — точки в прогоне нет" + ) + + +@pytest.mark.asyncio +async def test_serp_ok_done_exit_carries_checkpoint() -> None: + """Ранний done-выход #1950 («SERP собран, detail заблокирован») — тоже. + + Именно этим выходом кончается типичный прод-прогон, и он происходит РАНЬШЕ + единственной строки, которая писала точку. + """ + from scraper_kit.orchestration import pipeline as pl + + db = await _run( + raise_on=(ANCHOR_B[0], ANCHOR_B[1]), + exc=pl.AvitoBlockedError, + saved=(1, 0), # SERP intake > 0 → ветка ставит 'done', а не 'banned' + lots_per_anchor=1, + ) + + last = db.writes[-1] + assert "enrichment_abort_note" in last, "сработала не та ветка выхода" + assert last.get("done_buckets") == ["ekb-center"], ( + "ранний done-выход потерял якорь, пройденный до блокировки" + ) + + +@pytest.mark.asyncio +async def test_timed_out_anchor_is_not_checkpointed() -> None: + """Якорь, умерший по таймауту, не считается пройденным. + + Иначе резюм пропустит его навсегда, и это будет незаметно: прогон + завершается штатно, просто часть города не собирается никогда. + """ + db = await _run(raise_on=(ANCHOR_A[0], ANCHOR_A[1]), exc=TimeoutError) + + ckpt = db.writes[-1].get("done_buckets") + assert "ekb-center" not in ckpt, "якорь-таймаут попал в чекпоинт" + assert "ekb-south" in ckpt, "исправный якорь не зафиксирован" + + +@pytest.mark.asyncio +async def test_drain_exit_is_distinguishable_from_full_done() -> None: + """SIGTERM-дрейн помечен interrupted=1; полный обход — нет.""" + drained = await _run(shutdown_after_first=True) + full = await _run() + + assert drained.writes[-1].get("interrupted") == 1, ( + "оборванный дрейном прогон неотличим от полного обхода" + ) + assert drained.writes[-1].get("done_buckets") == ["ekb-center"] + assert "interrupted" not in full.writes[-1], "полный обход помечен как оборванный" + + +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_drained_done_is_resumable_but_clean_done_is_not() -> None: + """Метка дрейна доходит до решения о резюме — иначе она диагностика ради себя.""" + from scraper_kit.orchestration.scheduler import _resume_decision + + ckpt = {"done_buckets": ["ekb-center"], "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, "полный обход подхватывать нечего" + 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 71599c87..2cbacf9c 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 @@ -1183,6 +1183,18 @@ async def run_avito_city_sweep( ) _done_anchors: set[str] = set(_skip_anchors) + def _ckpt(**extra: Any) -> dict[str, Any]: + """Счётчики прогона ВМЕСТЕ с чекпоинтом — payload любого выхода (#3319). + + До этого точку писала ровно одна строка внутри цикла якорей, а все + финализаторы (mark_done/mark_banned/mark_failed, включая ранний выход + #1950 «SERP OK, detail заблокирован») отдавали голый `counters.to_dict()`. + Точка держалась исключительно на jsonb-мерже в runs.py — на свойстве + ЧУЖОГО модуля, которого этот файл ничем не проверяет; выход, случившийся + раньше первой записи (или писатель без мержа), терял её молча. + """ + return {**counters.to_dict(), "done_buckets": sorted(_done_anchors), **extra} + _loc = get_city_location(city_slug) # #262 wave 2: avito_slug у CityLocation Optional — не у каждого известного города # он подтверждён (403/429 на исчерпанном пуле при проверке, либо omonym-коллизия). @@ -1287,7 +1299,7 @@ async def run_avito_city_sweep( len(_anchors), name, ) - runs.update_heartbeat(db, run_id, counters.to_dict()) + runs.update_heartbeat(db, run_id, _ckpt()) return counters elif shutdown_requested(): # Кооперативный SIGTERM-drain (#1182 Phase 3a): останавливаемся @@ -1301,8 +1313,13 @@ async def run_avito_city_sweep( len(_anchors), name, ) - runs.update_heartbeat(db, run_id, counters.to_dict()) - runs.mark_done(db, run_id, counters.to_dict()) + # #3319: 'done' с counters.interrupted=1 — НЕ полный обход + # (та же метка, что у rosreestr_dkp-дрейна). Без неё оборванный + # деплоем прогон неотличим от честно обошедшего все якоря: + # статус тот же, счётчики частичные, а резюм его не берёт. + # Читатели статуса не трогаем — 'done' остаётся 'done'. + runs.update_heartbeat(db, run_id, _ckpt(interrupted=1)) + runs.mark_done(db, run_id, _ckpt(interrupted=1)) return counters logger.info( @@ -1762,6 +1779,12 @@ async def run_avito_city_sweep( _avito_anchor_timeout, ) counters.errors_count += 1 + # #3319: тот же инвариант, что у generic-except ниже. Якорь, + # умерший по таймауту, ПРОЙДЕН НЕ БЫЛ: SERP мог успеть, а + # detail/houses — нет, и какая именно часть осталась несобранной, + # здесь неизвестно. Считать его пройденным значит, что резюм + # пропустит его навсегда — молча, при штатно завершившемся прогоне. + _anchor_ok = False except (AvitoBlockedError, AvitoRateLimitedError) as e: logger.error( "city-sweep run_id=%d ABORT at anchor #%d/%d (%s) — blocked: %s", @@ -1773,7 +1796,7 @@ async def run_avito_city_sweep( ) counters.errors_count += 1 counters.anchors_done = idx - runs.update_heartbeat(db, run_id, counters.to_dict()) + runs.update_heartbeat(db, run_id, _ckpt()) # #1950: если SERP уже собрал лоты и заблокировало только detail/houses, # ставим 'done' (не 'banned') — partial intake сохранён. # За флагом avito_serp_ok_not_banned (default True). @@ -1794,14 +1817,14 @@ async def run_avito_city_sweep( runs.mark_done( db, run_id, - {**counters.to_dict(), "enrichment_abort_note": _note}, # type: ignore[arg-type] + _ckpt(enrichment_abort_note=_note), # type: ignore[arg-type] ) else: runs.mark_banned( db, run_id, str(e), - counters.to_dict(), + _ckpt(), ban_kind=ban_kind_of_exception(e), ) return counters @@ -1819,9 +1842,7 @@ async def run_avito_city_sweep( # завершится штатно. Тот же инвариант, что у combo в yandex-свипе. if _anchor_ok: _done_anchors.add(name) - runs.update_heartbeat( - db, run_id, {**counters.to_dict(), "done_buckets": sorted(_done_anchors)} - ) + runs.update_heartbeat(db, run_id, _ckpt()) # ── IMV-фаза: финальный обход тронутых домов ────────── if enrich_imv and all_touched_house_ids: @@ -1831,7 +1852,7 @@ async def run_avito_city_sweep( run_id, len(all_touched_house_ids), ) - runs.mark_done(db, run_id, counters.to_dict()) + runs.mark_done(db, run_id, _ckpt()) return counters elif shutdown_requested(): # SIGTERM-drain до IMV-фазы: финализируем без дорогой IMV-оценки. @@ -1841,7 +1862,10 @@ async def run_avito_city_sweep( run_id, len(all_touched_house_ids), ) - runs.mark_done(db, run_id, counters.to_dict()) + # #3319: тоже дрейн, но якоря пройдены ВСЕ — резюмить нечего + # (пропустил бы весь список и собрал ноль), поэтому метка + # диагностическая, а не резюм-флаг `interrupted`. + runs.mark_done(db, run_id, _ckpt(imv_phase_drained=1)) return counters logger.info( @@ -1855,7 +1879,7 @@ async def run_avito_city_sweep( # без update_heartbeat, и reap_zombies помечает живой run # 'zombie' → последующий mark_done становится no-op (дубль-sweep). def _imv_heartbeat() -> None: - runs.update_heartbeat(db, run_id, counters.to_dict()) + runs.update_heartbeat(db, run_id, _ckpt()) imv_result = await enrichment.process_houses_imv_batch( db, @@ -1867,7 +1891,7 @@ async def run_avito_city_sweep( counters.imv_enriched += imv_result.saved counters.imv_failed += imv_result.errors counters.errors_count += imv_result.errors - runs.update_heartbeat(db, run_id, counters.to_dict()) + runs.update_heartbeat(db, run_id, _ckpt()) logger.info( "city-sweep run_id=%d: IMV phase done — attempted=%d enriched=%d failed=%d", run_id, @@ -1892,9 +1916,9 @@ async def run_avito_city_sweep( db.rollback() except Exception: pass - runs.update_heartbeat(db, run_id, counters.to_dict()) + runs.update_heartbeat(db, run_id, _ckpt()) - runs.mark_done(db, run_id, counters.to_dict()) + runs.mark_done(db, run_id, _ckpt()) logger.info( "city-sweep run_id=%d done: anchors=%d/%d lots=%d (ins=%d/upd=%d) " "houses=%d/%d detail=%d/%d imv=%d/%d errors=%d", @@ -1916,7 +1940,7 @@ async def run_avito_city_sweep( except Exception as exc: logger.exception("city-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), _ckpt()) raise diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py index 0892d454..d49712c7 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py @@ -638,7 +638,14 @@ def _resume_decision(row: Any) -> tuple[int | None, dict[str, Any]]: } _boot_reaped_zombie = row.prev_status == "zombie" and prev_counters.get("boot_reaped") is True - if row.prev_status not in _RESUME_STATUSES and not _boot_reaped_zombie: + # 'done' с counters.interrupted=1 — SIGTERM-drain (#3319): статус штатный, но обход + # оборван на границе корзины, часть дерева не собрана. Сам статус в _RESUME_STATUSES + # не добавлен НАМЕРЕННО: чистое 'done' — полный проход, резюмить у него нечего, а + # подхват такой точки означал бы, что источник больше никогда не обходится целиком. + # Метка — та же, что у rosreestr_dkp-дрейна (app/services/scheduler.py), поэтому ни + # один читатель статуса не меняется. + _drained_done = row.prev_status == "done" and bool(prev_counters.get("interrupted")) + if row.prev_status not in _RESUME_STATUSES and not _boot_reaped_zombie and not _drained_done: verdict["resume_reason"] = f"status_{row.prev_status}" elif not row.same_params: verdict["resume_reason"] = "params_changed"