All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / changes (pull_request) Successful in 10s
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 5m13s
После #3319 резюм подхватывает 'done'-прогоны только с counters.interrupted (_drained_done в scheduler._resume_decision), но ставил метку ровно один avito_city_sweep. У yandex, cian и newbuilding SIGTERM-drain финализировался чистым 'done' с частичными счётчиками: оборванный деплоем обход неотличим от полного и из резюма выпадал, хотя чекпоинт done_buckets есть у всех трёх (combo-метки / имена якорей / номера страниц). done_buckets в дрейн-payload не добавляю: heartbeat мержит jsonb, уже записанные единицы обхода переживают финализатор, а пустой список у multi-anchor yandex затёр бы унаследованный при claim чекпоинт.
212 lines
8.4 KiB
Python
212 lines
8.4 KiB
Python
"""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"
|