"""#2725: сигнал живости слался один раз — до батча, — и живые прогоны reap'ились. Что было. `_execute_cian_backfill` дёргал `update_heartbeat` ровно один раз, ДО `backfill_cian_history()`, а сам батч (до 100 объявлений + 37 домов, каждое — fetch через браузер + пауза ~5 с) heartbeat не трогал. `reap_zombies` меряет именно `heartbeat_at` с порогом ZOMBIE_THRESHOLD_HOURS = 6 ч → прогон помечался 'zombie' строго на 6-м часу независимо от того, работает он или висит. Прод-замер 2026-08-06: 6 прогонов `cian_history_backfill` со статусом 'zombie', у всех шести сдвиг heartbeat 16-32 мс (= единственный стартовый вызов) и финал ровно на started_at + 6.00 ч. Живыми они при этом были: внутри окна пятерых писались строки offer_price_history с source='cian' (98/523/38/60/83 — плановый писатель этих строк только этот батч), у прогона 304 последняя строка легла через 5.40 ч после старта. Штатная длительность источника доходит до 5.06 ч (прогон 346, counters.duration_sec 18230) — то есть источник ходит вплотную к порогу. Почему чинится сигнал, а не критерий: пометка 'zombie' снимает running-блокировку источника (`has_running_run` видит только status='running'), и без неё зависший прогон запер бы источник навсегда. Плюс `mark_done` апдейтит `WHERE status='running'`, поэтому после ложной пометки собственный финал прогона — no-op (отсюда нулевые counters у всех шести строк). Фальсификация: на старом коде тест 2 падает — планировщик не передавал `on_progress`, батч heartbeat не двигал, и к концу 7-часовой работы возраст сигнала = 7 ч > 6 ч. """ from __future__ import annotations import os from datetime import UTC, datetime, timedelta from types import SimpleNamespace from typing import Any from unittest.mock import AsyncMock, MagicMock, patch os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") from scraper_kit.orchestration.scheduler import ZOMBIE_THRESHOLD_HOURS from app.services import scheduler as sched_mod from app.tasks import cian_history_backfill def _would_be_reaped(heartbeat_at: datetime, now: datetime) -> bool: """Критерий reap_zombies дословно: heartbeat старше порога → 'zombie'.""" return heartbeat_at < now - timedelta(hours=ZOMBIE_THRESHOLD_HOURS) class _FakeBrowserFetcher: def __init__(self, **kwargs: Any) -> None: pass async def __aenter__(self) -> _FakeBrowserFetcher: return self async def __aexit__(self, *_: object) -> None: return None # ── 1. Батч сообщает о продвижении на каждой сущности ──────────────────────── async def test_batch_reports_progress_per_entity() -> None: db = MagicMock() db.execute.return_value.mappings.return_value.all.side_effect = [ [ {"id": 1, "source_url": "https://cian.ru/1"}, {"id": 2, "source_url": "https://cian.ru/2"}, ], [{"id": 10, "cian_zhk_url": "https://cian.ru/zhk-10"}], ] seen: list[tuple[int, int]] = [] with ( patch.object(cian_history_backfill, "BrowserFetcher", _FakeBrowserFetcher), patch.object( cian_history_backfill, "fetch_detail", AsyncMock(return_value=SimpleNamespace(price_changes=[])), ), patch.object(cian_history_backfill, "save_detail_enrichment", MagicMock()), patch( "scraper_kit.providers.cian.newbuilding.fetch_newbuilding", AsyncMock(return_value=SimpleNamespace()), ), patch("scraper_kit.providers.cian.newbuilding.save_newbuilding_enrichment", MagicMock()), patch("asyncio.sleep", new_callable=AsyncMock), ): result = await cian_history_backfill.backfill_cian_history( db, do_listings=True, do_houses=True, do_valuations=False, on_progress=lambda r: seen.append((r.listings_processed, r.houses_processed)), ) # По одному сигналу на каждую сущность обоих блоков, счётчики растут. assert seen == [(1, 0), (2, 0), (2, 1)] assert result.listings_processed == 2 assert result.houses_processed == 1 # ── 2. Долгий прогон с продвигающимся heartbeat не помечается зависшим ─────── async def test_long_run_with_advancing_heartbeat_is_not_reaped() -> None: """7 часов работы, час на сущность: планировщик обязан двигать heartbeat.""" t0 = datetime(2026, 6, 26, 4, 47, tzinfo=UTC) clock = SimpleNamespace(now=t0) beats: list[datetime] = [] async def _fake_batch(db: Any, **kwargs: Any) -> Any: on_progress = kwargs.get("on_progress") result = cian_history_backfill.CianBackfillResult() for _ in range(7): # 7 сущностей по часу — дольше 6-часового порога clock.now += timedelta(hours=1) result.listings_processed += 1 if on_progress is not None: on_progress(result) result.duration_sec = 7 * 3600 return result fake_runs = SimpleNamespace( update_heartbeat=lambda db, run_id, counters: beats.append(clock.now), mark_done=MagicMock(), mark_failed=MagicMock(), ) with ( patch.object(sched_mod, "runs_mod", fake_runs), patch.object(cian_history_backfill, "backfill_cian_history", _fake_batch), ): await sched_mod._execute_cian_backfill(MagicMock(), run_id=1, params={}) assert len(beats) == 8, "стартовый сигнал + по одному на сущность" reaped = _would_be_reaped(beats[-1], clock.now) assert not reaped, "живой прогон с продвигающимся heartbeat не должен reap'иться" assert fake_runs.mark_done.called # ── 3. Прогон без продвижения — помечается (контроль критерия) ─────────────── async def test_long_run_without_advancing_heartbeat_is_reaped() -> None: """Тот же прогон, но батч сигнала не шлёт — критерий обязан сработать.""" t0 = datetime(2026, 6, 26, 4, 47, tzinfo=UTC) clock = SimpleNamespace(now=t0) beats: list[datetime] = [] async def _mute_batch(db: Any, **kwargs: Any) -> Any: clock.now += timedelta(hours=7) # работает, но молча return cian_history_backfill.CianBackfillResult() fake_runs = SimpleNamespace( update_heartbeat=lambda db, run_id, counters: beats.append(clock.now), mark_done=MagicMock(), mark_failed=MagicMock(), ) with ( patch.object(sched_mod, "runs_mod", fake_runs), patch.object(cian_history_backfill, "backfill_cian_history", _mute_batch), ): await sched_mod._execute_cian_backfill(MagicMock(), run_id=1, params={}) assert beats == [t0], "единственный сигнал — стартовый" assert _would_be_reaped(beats[-1], clock.now) # ── 4. Сбой heartbeat не роняет уже идущую работу ──────────────────────────── async def test_heartbeat_failure_does_not_abort_the_batch() -> None: processed: list[int] = [] calls = {"n": 0} def _flaky_heartbeat(db: Any, run_id: int, counters: dict[str, int]) -> None: calls["n"] += 1 if calls["n"] > 1: # стартовый прошёл, дальше БД отвалилась raise Exception("DB gone") async def _fake_batch(db: Any, **kwargs: Any) -> Any: on_progress = kwargs["on_progress"] result = cian_history_backfill.CianBackfillResult() for _ in range(3): result.listings_processed += 1 on_progress(result) # обязан проглотить исключение внутри себя processed.append(result.listings_processed) return result fake_runs = SimpleNamespace( update_heartbeat=_flaky_heartbeat, mark_done=MagicMock(), mark_failed=MagicMock(), ) with ( patch.object(sched_mod, "runs_mod", fake_runs), patch.object(cian_history_backfill, "backfill_cian_history", _fake_batch), ): await sched_mod._execute_cian_backfill(MagicMock(), run_id=1, params={}) assert processed == [1, 2, 3] assert fake_runs.mark_done.called # ── 5. Тот же дефект у newbuilding_enrich — сигнал прокинут ────────────────── async def test_newbuilding_enrich_passes_progress_callback() -> None: from app.tasks import newbuilding_enrich_backfill as nb beats: list[dict[str, int]] = [] async def _fake_backfill(db: Any, **kwargs: Any) -> Any: on_progress = kwargs.get("on_progress") assert on_progress is not None, "планировщик обязан прокинуть сигнал живости" result = nb.NewbuildingEnrichBackfillResult() result.processed += 1 on_progress(result) return result fake_runs = SimpleNamespace( update_heartbeat=lambda db, run_id, counters: beats.append(counters), mark_done=MagicMock(), mark_failed=MagicMock(), # Финал прогона ушёл в общий mark_backfill_finished (#2767) — без него подмена # runs_mod роняет AttributeError и тест меряет не то, что проверяет. mark_backfill_finished=MagicMock(), ) with ( patch.object(nb, "runs_mod", fake_runs), patch.object(nb, "backfill_newbuilding_enrichment", _fake_backfill), ): await nb.run_newbuilding_enrich(MagicMock(), run_id=1, params={}) assert len(beats) == 2, "стартовый сигнал + сигнал из середины цикла" assert beats[-1]["processed"] == 1