All checks were successful
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m2s
Deploy Trade-In / build-backend (push) Successful in 1m34s
Deploy Trade-In / deploy (push) Successful in 2m11s
228 lines
11 KiB
Python
228 lines
11 KiB
Python
"""#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
|