fix(scraper-kit): чекпоинт avito_city_sweep доживает до финализатора — 0/67 прогонов писали done_buckets #3330
3 changed files with 279 additions and 17 deletions
225
tradein-mvp/backend/tests/test_3319_citysweep_checkpoint.py
Normal file
225
tradein-mvp/backend/tests/test_3319_citysweep_checkpoint.py
Normal 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"
|
||||
|
|
@ -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
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue