fix(scraper-kit): метка дрейна interrupted=1 во всех city-свипах — yandex/cian/newbuilding больше не объявляют недоделанный обход готовым #3346
2 changed files with 230 additions and 6 deletions
212
tradein-mvp/backend/tests/test_3333_drain_mark_all_sweeps.py
Normal file
212
tradein-mvp/backend/tests/test_3333_drain_mark_all_sweeps.py
Normal 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"
|
||||
|
|
@ -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(
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue