gendesign/tradein-mvp/backend/tests/test_2725_heartbeat_in_batch.py
bot-backend 5046ac7b4e
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
fix(tradein/cian): показать, ЧТО пришло вместо состояния ЖК-страницы, и честно закрывать нулевой прогон (#2767) (#2768)
2026-08-07 08:06:25 +00:00

228 lines
11 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""#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