fix(scraper-kit): чекпоинт avito_city_sweep доживает до финализатора — 0/67 прогонов писали done_buckets #3330

Merged
bot-backend merged 2 commits from fix/3319-citysweep-checkpoint into main 2026-09-05 17:36:51 +00:00
3 changed files with 279 additions and 17 deletions

View file

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

View file

@ -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

View file

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