fix(scraper): пометить SIGTERM-дрейн interrupted=1 у full-load'ов и domclick-свипа
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 12s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m5s

После #3330/#3346 резюм берёт 'done' только с counters.interrupted=1, но четыре
функции с живым чекпоинтом done_buckets финализировали дрейн чистым 'done':
cian/yandex/avito full-load (ветка RuntimeError("shutdown")) и domclick city
sweep (дрейн до SERP). Оборванный обход был неотличим от полного, а собранные
бакеты никто не подхватывал — у avito это ещё и бан-бюджет (#3315).

done_buckets в domclick-payload не добавляем: чекпоинт унаследован при claim
(#3074), jsonb-мерж его сохраняет.

Closes #3355
This commit is contained in:
bot-backend 2026-09-05 23:59:22 +05:00
parent 63dbc209b2
commit f6cf948358
2 changed files with 302 additions and 5 deletions

View file

@ -0,0 +1,269 @@
"""SIGTERM-дрейн помечен `interrupted=1` и у full-load'ов с domclick-свипом (#3355).
#3330/#3346 закрыли city-свипы, но четыре функции с ЖИВЫМ чекпоинтом `done_buckets`
продолжали финализировать дрейн чистым `done`: cian/yandex/avito full-load (ветка
`RuntimeError("shutdown")`) и domclick city sweep (дрейн до SERP-фазы). Резюм
(`_drained_done` в scheduler._resume_decision) такой прогон не берёт чекпоинт
пишется в никуда.
Кто чекпоинт ЧИТАЕТ (метка не диагностическая):
* cian_full_load `skip_set` из `done_buckets` (ключ "room:lo:hi"),
планировщик зовёт с `resume_run_id=_pick_resume(...)`;
* avito_full_load то же, ДВА job'а (обычный и exhaustive, #3315);
* domclick_city_sweep `skip_buckets` из `done_buckets` (ROOM_BUCKETS), тоже
`_pick_resume`;
* yandex_full_load читает `done_buckets` так же, НО у него нет job'а в
scheduler.py: единственный вход с `resume_run_id`
админский POST. Метка там честно работает только при
ручном подхвате; автоматического потребителя нет.
Статус финализации нулевого дрейна фиксируется тестом отдельно: honest-status
гейты (runs.mark_done _sweep_run_did_nothing/_phase_totally_failed/
_failed_ratio_too_high) могут увести прогон в 'failed'. Это резюмируемо
('failed' _RESUME_STATUSES), но статус факт, а не догадка.
Анти-цикл («чистый done не резюмируется») уже закрыт в
test_3333_drain_mark_all_sweeps.py здесь не дублируется.
"""
from __future__ import annotations
import os
# Settings собирается автофикстурой conftest'а и требует database_url — как в
# test_3333_drain_mark_all_sweeps.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
class _FakeDb:
"""Каждый UPDATE с counters запоминается вместе с текстом SQL (для статуса)."""
def __init__(self) -> None:
self.writes: list[tuple[str, dict[str, Any]]] = []
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
if params and "counters" in params:
self.writes.append((str(stmt), json.loads(params["counters"])))
return MagicMock()
def commit(self) -> None: ...
def rollback(self) -> None: ...
def _final(db: _FakeDb) -> tuple[str, dict[str, Any]]:
"""Последняя запись counters = финализатор."""
assert db.writes, "дрейн не оставил ни одной записи counters"
sql, counters = db.writes[-1]
if "status = 'done'" in sql:
return "done", counters
if "status = 'failed'" in sql:
return "failed", counters
return "heartbeat", counters
class _DrainAtFirstBucket:
"""Любой fetch-метод сразу зовёт on_bucket — настоящий `_on_bucket` и роняет дрейн."""
def __init__(self, *_a: Any, **_kw: Any) -> None:
self._browser = None
self._cffi = None
self.request_delay_sec = 0.0
async def __aenter__(self) -> _DrainAtFirstBucket:
return self
async def __aexit__(self, *_e: Any) -> None:
return None
def __getattr__(self, name: str) -> Any:
async def _fetch(*_a: Any, **kw: Any) -> Any:
on_bucket = kw.get("on_bucket")
assert on_bucket is not None, f"{name} вызван без on_bucket — тест мимо дрейна"
on_bucket("2:0:5000000", [])
raise AssertionError("_on_bucket не оборвал прогон при shutdown_requested()")
return _fetch
class _NeverCalledScraper:
"""Дрейн срабатывает ДО скрапера: любой запрос тут — сломанный порядок проверок."""
def __init__(self, *_a: Any, **_kw: Any) -> None: ...
async def __aenter__(self) -> _NeverCalledScraper:
return self
async def __aexit__(self, *_e: Any) -> None:
return None
def __getattr__(self, name: str) -> Any:
async def _boom(*_a: Any, **_kw: Any) -> Any:
raise AssertionError(f"скрапер вызван при дрейне: {name}")
return _boom
def _config() -> types.SimpleNamespace:
return types.SimpleNamespace(
scraper_fetch_mode="cffi",
scraper_proxy_url=None,
scraper_skip_seen_today=False,
use_proxy_pool_browser=False,
browser_http_endpoint=None,
environment="test",
avito_serp_ok_not_banned=True,
cian_full_load_per_fetch_timeout_s=0.0,
)
@pytest.mark.asyncio
async def test_cian_full_load_drain_is_marked_interrupted() -> None:
"""cian full-load: дрейн на границе бакета → counters.interrupted == 1."""
from scraper_kit.orchestration import pipeline as pl
db = _FakeDb()
with (
patch.object(pl, "CianScraper", _DrainAtFirstBucket),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
await pl.run_cian_full_load(
db, # type: ignore[arg-type]
run_id=3355,
config=_config(),
matcher=MagicMock(),
request_delay_sec=0.0,
shutdown_requested=lambda: True,
)
_status, counters = _final(db)
assert counters.get("interrupted") == 1, (
"cian full-load: оборванный дрейном прогон неотличим от полного обхода — "
"резюм его чекпоинт не подхватит"
)
@pytest.mark.asyncio
async def test_yandex_full_load_drain_is_marked_interrupted() -> None:
"""yandex full-load: дрейн на границе бакета → counters.interrupted == 1."""
from scraper_kit.orchestration import pipeline as pl
db = _FakeDb()
with (
patch.object(pl, "YandexRealtyScraper", _DrainAtFirstBucket),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
await pl.run_yandex_full_load(
db, # type: ignore[arg-type]
run_id=3355,
config=_config(),
matcher=MagicMock(),
enrichment=MagicMock(),
request_delay_sec=0.0,
shutdown_requested=lambda: True,
)
_status, counters = _final(db)
assert counters.get("interrupted") == 1, (
"yandex full-load: оборванный дрейном прогон неотличим от полного обхода"
)
@pytest.mark.asyncio
async def test_avito_full_load_drain_is_marked_interrupted() -> None:
"""avito full-load (exhaustive, #3315): дрейн на границе бакета → interrupted == 1."""
from scraper_kit.orchestration import pipeline as pl
db = _FakeDb()
with (
patch.object(pl, "AvitoScraper", _DrainAtFirstBucket),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
await pl.run_avito_full_load(
db, # type: ignore[arg-type]
run_id=3355,
config=_config(),
matcher=MagicMock(),
request_delay_sec=0.0,
shutdown_requested=lambda: True,
)
_status, counters = _final(db)
assert counters.get("interrupted") == 1, (
"avito full-load: дрейн без метки — бан-бюджет тратится на уже собранные бакеты"
)
@pytest.mark.asyncio
async def test_domclick_sweep_drain_is_marked_interrupted() -> None:
"""domclick city sweep: дрейн до SERP-фазы → counters.interrupted == 1."""
from scraper_kit.orchestration import pipeline as pl
db = _FakeDb()
with (
patch.object(pl, "DomClickScraper", _NeverCalledScraper),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
await pl.run_domclick_city_sweep(
db, # type: ignore[arg-type]
run_id=3355,
config=_config(),
matcher=MagicMock(),
request_delay_sec=0.0,
shutdown_requested=lambda: True,
)
_status, counters = _final(db)
assert counters.get("interrupted") == 1, (
"domclick: оборванный дрейном прогон неотличим от полного обхода"
)
@pytest.mark.asyncio
async def test_zero_drain_status_is_resumable() -> None:
"""Минор ревью #3346: какой СТАТУС даёт нулевой дрейн — и берёт ли его резюм.
honest-status гейты в runs.mark_done (#2625/#2700/honest-run-status) вправе
увести прогон в 'failed'. Тест фиксирует ФАКТ, а не ожидание: обе ветки
резюмируемы ('failed' через _RESUME_STATUSES, 'done' через метку), но
третьего исхода быть не должно.
"""
from scraper_kit.orchestration import pipeline as pl
from scraper_kit.orchestration.scheduler import _RESUME_STATUSES, _resume_decision
db = _FakeDb()
with (
patch.object(pl, "AvitoScraper", _DrainAtFirstBucket),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
await pl.run_avito_full_load(
db, # type: ignore[arg-type]
run_id=3355,
config=_config(),
matcher=MagicMock(),
request_delay_sec=0.0,
shutdown_requested=lambda: True,
)
status, counters = _final(db)
assert status == "done", f"нулевой дрейн финализирован как {status!r} (было 'done')"
prev = types.SimpleNamespace(
prev_id=4707,
prev_status=status,
prev_counters={**counters, "done_buckets": ["2:0:5000000"], "resume_chain": 0},
same_params=True,
age_h=2.0,
interval_days="7",
)
resume_from, verdict = _resume_decision(prev)
assert resume_from == 4707, (
f"дрейн со статусом {status!r} не подхвачен резюмом: {verdict} "
f"(_RESUME_STATUSES={sorted(_RESUME_STATUSES)})"
)

View file

@ -3755,7 +3755,15 @@ async def run_cian_full_load(
counters.saved_inserted,
counters.saved_updated,
)
runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
# #3355: та же метка дрейна, что у city-свипов (#3330/#3333). Без неё
# оборванный дрейном обход неотличим от полного: статус тот же 'done',
# а `_drained_done` в scheduler._resume_decision смотрит именно на неё —
# и done_buckets этого прогона не подхватывал никто.
runs.mark_done(
db,
run_id,
{**counters.to_dict(), "done_buckets": sorted(done), "interrupted": 1},
)
return counters
logger.exception("cian-full-load run_id=%d: fatal error", run_id)
runs.mark_failed(db, run_id, str(exc), counters.to_dict())
@ -3997,7 +4005,15 @@ async def run_yandex_full_load(
counters.saved_inserted,
counters.saved_updated,
)
runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
# #3355: та же метка дрейна, что у city-свипов (#3330/#3333). Без неё
# оборванный дрейном обход неотличим от полного: статус тот же 'done',
# а `_drained_done` в scheduler._resume_decision смотрит именно на неё —
# и done_buckets этого прогона не подхватывал никто.
runs.mark_done(
db,
run_id,
{**counters.to_dict(), "done_buckets": sorted(done), "interrupted": 1},
)
return counters
logger.exception("yandex-full-load run_id=%d: fatal error", run_id)
runs.mark_failed(db, run_id, str(exc), counters.to_dict())
@ -4263,7 +4279,15 @@ async def run_avito_full_load(
counters.saved_inserted,
counters.saved_updated,
)
runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
# #3355: та же метка дрейна, что у city-свипов (#3330/#3333). Без неё
# оборванный дрейном обход неотличим от полного: статус тот же 'done',
# а `_drained_done` в scheduler._resume_decision смотрит именно на неё —
# и done_buckets этого прогона не подхватывал никто.
runs.mark_done(
db,
run_id,
{**counters.to_dict(), "done_buckets": sorted(done), "interrupted": 1},
)
return counters
logger.exception("avito-full-load run_id=%d: fatal error", run_id)
runs.mark_failed(db, run_id, str(exc), counters.to_dict())
@ -4385,8 +4409,12 @@ async def run_domclick_city_sweep(
# SIGTERM-drain до SERP-фазы: финализируем run как mark_done(partial),
# минуя honest-status (это чистый drain, не QRATOR-блок).
logger.info("domclick-sweep run_id=%d: SIGTERM-drain — stopping before SERP", run_id)
runs.update_heartbeat(db, run_id, counters.to_dict())
runs.mark_done(db, run_id, counters.to_dict())
# #3355: метка дрейна (см. city-свипы #3330/#3333). done_buckets тут не
# пишем — унаследованный при claim (#3074) чекпоинт лежит в counters, и
# jsonb-мерж heartbeat/финализатора его сохраняет.
_drain = {**counters.to_dict(), "interrupted": 1}
runs.update_heartbeat(db, run_id, _drain)
runs.mark_done(db, run_id, _drain)
return counters
logger.info(