fix(scraper-kit): метка дрейна interrupted=1 во всех city-свипах — yandex/cian/newbuilding больше не объявляют недоделанный обход готовым #3346

Merged
bot-backend merged 1 commit from fix/3333-drain-mark-all-sweeps into main 2026-09-05 18:03:05 +00:00
2 changed files with 230 additions and 6 deletions

View file

@ -0,0 +1,212 @@
"""SIGTERM-дрейн помечен `interrupted=1` во ВСЕХ свипах, не только у avito (#3333).
После #3319 резюм подхватывает 'done'-прогоны только с меткой `interrupted`
(`_drained_done` в scheduler._resume_decision), а ставил её ровно один
avito_city_sweep. У yandex (~2356 стр.), cian (~2960) и newbuilding (~2098)
дрейн финализировался чистым `done` с частичными счётчиками: недоделанный обход
объявлен полным, статус тот же, что у честного, и из резюма он выпадал.
Резюм есть у всех трёх (чекпоинт `done_buckets`: combo-метки у yandex, имена
якорей у cian, номера страниц у newbuilding), так что метка не диагностическая:
у каждого есть что подхватывать. Оговорка одна yandex подхватывает только при
единственном якоре (combo-ключ не содержит якоря); прод-режим ровно такой.
Анти-цикл (последний тест): метка не должна превратить ЛЮБОЙ 'done' в
резюмируемый иначе источник больше никогда не обходится целиком.
"""
from __future__ import annotations
import os
# Settings собирается автофикстурой conftest'а и требует database_url — как в
# test_3319_citysweep_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 _NeverCalledScraper:
"""Дрейн срабатывает ДО скрапера: любой запрос тут — сломанный порядок проверок."""
def __init__(self, *_a: Any, **_kw: Any) -> None:
self._browser = None
self._cffi = None
self.state_extraction_attempts = 1
self.state_extraction_failures = 0
self.request_delay_sec = 0.0
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,
use_proxy_pool_browser=False,
browser_http_endpoint=None,
environment="test",
avito_serp_ok_not_banned=True,
)
@pytest.mark.asyncio
async def test_yandex_drain_is_marked_interrupted() -> None:
"""yandex-sweep: дрейн на границе якоря → counters.interrupted == 1."""
from scraper_kit.orchestration import pipeline as pl
db = _FakeDb()
enrichment = MagicMock()
enrichment.record_yandex_price_history.return_value = 0
with (
patch.object(pl, "YandexRealtyScraper", _NeverCalledScraper),
patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
await pl.run_yandex_city_sweep(
db, # type: ignore[arg-type]
run_id=3333,
config=_config(),
matcher=MagicMock(),
enrichment=enrichment,
enrich_address=False,
shutdown_requested=lambda: True,
)
assert db.writes, "дрейн не оставил ни одной записи counters"
assert db.writes[-1].get("interrupted") == 1, (
"yandex: оборванный дрейном прогон неотличим от полного обхода — резюм его не возьмёт"
)
@pytest.mark.asyncio
async def test_cian_drain_is_marked_interrupted() -> None:
"""cian-sweep: дрейн на границе якоря → counters.interrupted == 1."""
from scraper_kit.orchestration import pipeline as pl
db = _FakeDb()
with (
patch.object(pl, "CianScraper", _NeverCalledScraper),
patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
await pl.run_cian_city_sweep(
db, # type: ignore[arg-type]
run_id=3333,
config=_config(),
matcher=MagicMock(),
anchors=[ANCHOR_A, ANCHOR_B],
enrich_houses=False,
detail_top_n=0,
request_delay_sec=0.0,
shutdown_requested=lambda: True,
)
assert db.writes, "дрейн не оставил ни одной записи counters"
assert db.writes[-1].get("interrupted") == 1, (
"cian: оборванный дрейном прогон неотличим от полного обхода — резюм его не возьмёт"
)
@pytest.mark.asyncio
async def test_newbuilding_drain_is_marked_interrupted() -> None:
"""nb-sweep: дрейн до SERP-фазы → counters.interrupted == 1."""
from scraper_kit.orchestration import pipeline as pl
db = _FakeDb()
with (
patch.object(pl, "AvitoScraper", _NeverCalledScraper),
patch.object(pl, "AsyncSession", _FakeAsyncSession),
patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
await pl.run_avito_newbuilding_sweep(
db, # type: ignore[arg-type]
run_id=3333,
config=_config(),
matcher=MagicMock(),
pages=4,
request_delay_sec=0.0,
shutdown_requested=lambda: True,
)
assert db.writes, "дрейн не оставил ни одной записи counters"
assert db.writes[-1].get("interrupted") == 1, (
"newbuilding: оборванный дрейном прогон неотличим от полного обхода — резюм его не возьмёт"
)
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_clean_done_is_not_resumed_but_marked_drain_is() -> None:
"""Анти-цикл: подхватывается ТОЛЬКО помеченный дрейн, чистый 'done' — нет.
Общий на все три провайдера: `_resume_decision` смотрит на counters, а не на
источник, поэтому разбор один. Если бы метка была не нужна для подхвата, все
три правки выше были бы записью в лог ради записи в лог.
"""
from scraper_kit.orchestration.scheduler import _resume_decision
ckpt = {"done_buckets": ["combo-1"], "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, "чистый 'done' — полный обход, подхватывать нечего"
assert verdict["resume_reason"] == "status_done"

View file

@ -2112,8 +2112,14 @@ async def run_avito_newbuilding_sweep(
logger.info(
"nb-sweep run_id=%d: SIGTERM-drain — stopping before SERP phase", run_id
)
runs.update_heartbeat(db, run_id, counters.to_dict())
runs.mark_done(db, run_id, counters.to_dict())
# #3333: та же метка дрейна, что у avito_city_sweep (_ckpt(interrupted=1)).
# Без неё оборванный деплоем обход неотличим от полного: статус 'done',
# счётчики частичные — и резюм (`_drained_done` в scheduler) его не берёт.
# done_buckets тут не пишем: heartbeat мержит jsonb, уже записанные
# страницы переживают финализатор.
_drain = {**counters.to_dict(), "interrupted": 1}
runs.update_heartbeat(db, run_id, _drain)
runs.mark_done(db, run_id, _drain)
return counters
# proxy_provider прокинут для консистентности (#2616) — не load-bearing,
@ -2370,8 +2376,11 @@ async def run_yandex_city_sweep(
len(_anchors),
name,
)
runs.update_heartbeat(db, run_id, counters.to_dict())
runs.mark_done(db, run_id, counters.to_dict())
# #3333: та же метка дрейна, что у avito_city_sweep (_ckpt(interrupted=1)) —
# см. там же. done_buckets пишет combo-heartbeat, jsonb-мерж их сохраняет.
_drain = {**counters.to_dict(), "interrupted": 1}
runs.update_heartbeat(db, run_id, _drain)
runs.mark_done(db, run_id, _drain)
return counters
logger.info(
@ -2985,8 +2994,11 @@ async def run_cian_city_sweep(
len(_anchors),
name,
)
runs.update_heartbeat(db, run_id, counters.to_dict())
runs.mark_done(db, run_id, counters.to_dict())
# #3333: та же метка дрейна, что у avito_city_sweep (_ckpt(interrupted=1)) —
# см. там же. done_buckets пишет end-of-anchor heartbeat, jsonb-мерж хранит.
_drain = {**counters.to_dict(), "interrupted": 1}
runs.update_heartbeat(db, run_id, _drain)
runs.mark_done(db, run_id, _drain)
return counters
logger.info(