Сбор МЕРА: упавший прогон не висит zombie 6 часов, провалившаяся фаза не попадает в чекпоинт, проверка отмены не держит транзакцию часами #3561

Merged
bot-backend merged 8 commits from fix/kit-orchestration into main 2026-09-17 09:23:44 +00:00
2 changed files with 39 additions and 6 deletions
Showing only changes of commit c7a03495e2 - Show all commits

View file

@ -3,7 +3,8 @@
Уже покрыто: avito таймаут якоря и дрейн после первого якоря (test_3319); дрейн ДО
первого якоря у yandex/cian/newbuilding (test_3333); отмена у domclick (test_3369).
Здесь остаток:
(а) таймаут якоря у cian и yandex, таймаут фазы у domclick;
(а) таймаут якоря у cian и yandex; у domclick снятая/упавшая фаза не отмечает
корзины, строки которых не сохранены;
(б) SIGTERM-дрейн ПОСЛЕ первого якоря у cian и yandex;
(в) пользовательская отмена (runs.is_cancelled=True) после первого якоря у avito,
cian и yandex.
@ -168,7 +169,16 @@ async def test_anchor_timeout_skips_anchor_and_continues(kind: str) -> None:
assert db.counters["errors_count"] >= 1
async def test_domclick_phase_timeout_keeps_completed_buckets_and_is_not_done() -> None:
@pytest.mark.parametrize("broken", ["fetch_timeout", "save_failed"])
async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> None:
"""Корзина без сохранённых строк не попадает в чекпоинт.
Лоты Домклика копятся в памяти и пишутся ОДНИМ save_listings после всех корзин.
Фаза, снятая watchdog'ом (или упавшая на save), не сохранила ничего — даже из
корзины 'st', которую скрейпер успел пройти. Отметить её пройденной значило бы,
что следующий прогон пропустит её через skip_buckets навсегда (миграция 308).
"""
class _Dc(_Scraper):
blocked = False
geo_filtered = fetch_errors = bucket_start_index = 0
@ -176,11 +186,22 @@ async def test_domclick_phase_timeout_keeps_completed_buckets_and_is_not_done()
completed_buckets = ["st"] # noqa: RUF012
async def fetch_city(self, **_kw: Any) -> list[Any]:
if broken == "fetch_timeout":
raise TimeoutError
return [object()]
saved: list[int] = []
def _save(_db: Any, lots: list[Any], **_kw: Any) -> tuple[int, int]:
if broken == "save_failed":
raise RuntimeError("save_listings упал")
saved.append(len(lots))
return len(lots), 0
db = _FakeDb()
with (
patch.object(pl, "DomClickScraper", _Dc),
patch.object(pl, "save_listings", _save),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
counters = await pl.run_domclick_city_sweep(
@ -192,8 +213,11 @@ async def test_domclick_phase_timeout_keeps_completed_buckets_and_is_not_done()
request_delay_sec=0.0,
)
assert saved == [], "тест не о том: строки сохранились"
assert counters.errors_count >= 1
assert db.counters["done_buckets"] == ["st"]
assert db.counters["done_buckets"] == [], (
f"корзины без единой сохранённой строки в чекпоинте: {db.counters['done_buckets']}"
)
assert db.status == "failed"

View file

@ -4992,10 +4992,13 @@ async def run_domclick_city_sweep(
)
lots: list[ScrapedLot] = []
# #2406: собранное сохраняется ОДНИМ save_listings после всех корзин, поэтому
# корзина считается пройденной только если фаза дошла до конца (см. чекпоинт ниже).
_saved = False
async def _domclick_phase() -> None:
"""Единственная citywide-фаза: fetch_city + save."""
nonlocal lots
nonlocal lots, _saved
async with DomClickScraper(
config,
proxy_provider=proxy_provider,
@ -5041,6 +5044,7 @@ async def run_domclick_city_sweep(
)
counters.lots_inserted += inserted
counters.lots_updated += updated
_saved = True
try:
await asyncio.wait_for(_domclick_phase(), timeout=_sweep_timeout)
@ -5094,6 +5098,11 @@ async def run_domclick_city_sweep(
# #3118: чекпоинт = унаследованное завершённое в этом прогоне.
# Пишем heartbeat'ом СЕЙЧАС (мерж jsonb) — финализаторы ключ не
# затирают, и оборванный болезнью финализации прогон его не теряет.
# #2406: только если лоты сохранены. Снятая watchdog'ом или упавшая фаза
# не дошла до save_listings — у её корзин в БД ноль строк, а отметка
# «пройдена» заставила бы следующий прогон пропустить их навсегда
# (механизм разобран в миграции 308, из-за него выключены свипы 77/50).
if _saved:
_checkpoint = sorted(skip_buckets | set(_s.completed_buckets))
runs.update_heartbeat(db, run_id, _payload())
counters.bucket_start_index = _s.bucket_start_index