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..a238119b --- /dev/null +++ b/tradein-mvp/backend/tests/test_3319_citysweep_checkpoint.py @@ -0,0 +1,225 @@ +"""Чекпоинт avito_city_sweep доживает до финализатора (#3319). + +Точку писала ровно одна строка — end-of-anchor heartbeat в конце итерации цикла +якорей. Финализаторы её не стирали (все писатели в runs.py мержат jsonb: +`counters || :counters`), дыра в другом: выходы, случившиеся РАНЬШЕ первой такой +записи, точки не оставляли вовсе — cancel/SIGTERM-дрейн на границе первого якоря +и ранний done #1950 («SERP собран, detail заблокирован») на якоре №1. Ими и +кончается типичный прод-прогон с `anchors_done: 1` из 5. + +Замер «0 из 67 прогонов за 60 дней несут done_buckets» тут НЕ доказательство: +строка записи появилась только 26.08.2026 (#3074) при такте avito 7 суток — +выборка почти целиком из эры, где механизма не существовало. + +Три инварианта, ради которых тест: + 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" + + # Прогон из эры до #3074: ключей нет вовсе — метка дрейна не должна менять + # вердикт «нечего подхватывать» на что-то другое. + _, verdict = _resume_decision(_prev_run({})) + assert verdict["resume_reason"] == "status_done" + _, verdict = _resume_decision(_prev_run({"interrupted": 1})) + assert verdict["resume_reason"] == "no_checkpoint" 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..654927ef 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,24 @@ async def run_avito_city_sweep( ) _done_anchors: set[str] = set(_skip_anchors) + def _ckpt(**extra: Any) -> dict[str, Any]: + """Счётчики прогона ВМЕСТЕ с чекпоинтом — payload любого выхода (#3319). + + До этого точку писала ровно одна строка — end-of-anchor heartbeat в конце + итерации цикла. Финализаторы её НЕ стирали: все четыре писателя в runs.py + мержат jsonb (`counters || :counters`), уже записанный ключ переживал и + mark_done, и mark_banned, и mark_failed. Закрывается другая дыра — выходы, + случившиеся РАНЬШЕ первой такой записи: cancel/SIGTERM-дрейн на границе + первого якоря и ранний done #1950 («SERP собран, detail заблокирован») на + якоре №1. Именно им и кончается типичный прод-прогон, у которого + `anchors_done: 1` из 5. + + Замер «0 из 67 прогонов за 60 дней несут done_buckets» сам по себе этого НЕ + доказывает: строка записи появилась только 26.08.2026 (#3074), а такт avito + — 7 суток, так что выборка почти целиком из эры, где механизма не было. + """ + 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 +1305,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 +1319,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 +1785,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 +1802,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 +1823,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 +1848,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 +1858,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 +1868,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 +1885,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 +1897,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 +1922,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 +1946,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"