All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 9s
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 4m51s
373 lines
16 KiB
Python
373 lines
16 KiB
Python
"""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,
|
||
prev_counters: dict[str, Any] | None = None,
|
||
detail_rows: list[dict[str, Any]] | None = None,
|
||
) -> None:
|
||
self.writes: list[tuple[str, dict[str, Any]]] = []
|
||
self._prev_counters = prev_counters
|
||
self._detail_rows = detail_rows or []
|
||
|
||
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()
|
||
if params and "rid" in params:
|
||
# SELECT counters предыдущего прогона (resume_run_id → чекпоинт).
|
||
row = types.SimpleNamespace(counters=self._prev_counters)
|
||
return MagicMock(**{"fetchone.return_value": row if self._prev_counters else None})
|
||
if params and "lim" in params:
|
||
# SELECT кандидатов detail-обогащения (cian full-load).
|
||
return MagicMock(**{"mappings.return_value.all.return_value": self._detail_rows})
|
||
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
|
||
# Данные-атрибуты объявляем явно: catch-all __getattr__ ниже отдаёт корутину
|
||
# на ЛЮБОЕ имя, и счётчик #3368 (читается в finally, ревью #3373) уехал бы
|
||
# в counters функцией — падало бы сериализацией, а не смыслом.
|
||
self.capped_buckets = 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_domclick_sweep_drain_keeps_inherited_checkpoint() -> None:
|
||
"""domclick: дрейн ОБЯЗАН переписать унаследованный чекпоинт в свой payload.
|
||
|
||
Мерж jsonb (#3390) тут не спасает: он сливает payload со СВОЕЙ строкой прогона, а та
|
||
создана пустой (`create_run`) — унаследованный чекпоинт лежит в counters ПРЕДЫДУЩЕГО
|
||
прогона, и scheduler при claim его не переносит. Без явного ключа дрейн-прогон
|
||
закрывался бы с пустым чекпоинтом, и резюм пересобирал бы все шесть корзин заново.
|
||
"""
|
||
from scraper_kit.orchestration import pipeline as pl
|
||
|
||
inherited = ["1:0:5000000", "2:0:5000000"]
|
||
db = _FakeDb(prev_counters={"done_buckets": inherited})
|
||
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,
|
||
resume_run_id=3354,
|
||
)
|
||
|
||
assert len(db.writes) >= 2, "дрейн обязан писать и heartbeat, и финализатор"
|
||
for label, (_sql, counters) in (("heartbeat", db.writes[-2]), ("финализатор", db.writes[-1])):
|
||
assert counters.get("done_buckets") == inherited, (
|
||
f"{label}: чекпоинт потерян ({counters.get('done_buckets')!r} вместо "
|
||
f"{inherited!r}) — резюм пересоберёт уже собранные корзины"
|
||
)
|
||
assert counters.get("interrupted") == 1, f"{label}: дрейн без метки"
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_cian_full_load_drain_in_detail_phase_is_marked_interrupted() -> None:
|
||
"""cian full-load: дрейн в detail-фазе — SERP целый, но прогон НЕ полный.
|
||
|
||
Ранний выход из detail-цикла делал `break` и уходил в обычный `mark_done`:
|
||
статус 'done' без метки, хотя обогащение обрезано на первой же записи.
|
||
"""
|
||
from scraper_kit.orchestration import pipeline as pl
|
||
|
||
class _SerpDoneScraper:
|
||
"""SERP отработал целиком (on_bucket не зовётся) — дрейн ловит detail-фаза."""
|
||
|
||
def __init__(self, *_a: Any, **_kw: Any) -> None:
|
||
self._browser = None
|
||
self.request_delay_sec = 0.0
|
||
self.last_dropped_nb = 0
|
||
self.state_extraction_attempts = 0
|
||
self.state_extraction_failures = 0
|
||
|
||
async def __aenter__(self) -> _SerpDoneScraper:
|
||
return self
|
||
|
||
async def __aexit__(self, *_e: Any) -> None:
|
||
return None
|
||
|
||
async def fetch_all_secondary(self, **_kw: Any) -> None:
|
||
return None
|
||
|
||
db = _FakeDb(detail_rows=[{"id": 1, "source_url": "https://cian.ru/1"}])
|
||
with (
|
||
patch.object(pl, "CianScraper", _SerpDoneScraper),
|
||
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,
|
||
enrich_detail=True,
|
||
detail_top_n=5,
|
||
shutdown_requested=lambda: True,
|
||
)
|
||
|
||
_status, counters = _final(db)
|
||
assert counters.get("interrupted") == 1, (
|
||
"cian full-load: дрейн в detail-фазе финализирован как полный прогон"
|
||
)
|
||
|
||
|
||
@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)})"
|
||
)
|