Ревью PR #3561: тест таймаута фазы закреплял как контракт done_buckets == ['st']. Лоты ДомКлика копятся в памяти и пишутся одним save_listings после всех корзин, поэтому снятая watchdog'ом (или упавшая на save) фаза не сохраняет ничего, а чекпоинт всё равно получал completed_buckets живого скрейпера — следующий прогон пропускал корзину навсегда (механизм миграции 308). Чекпоинт пополняется только если фаза дошла до конца (флаг _saved после save). Тест развёрнут: при таймауте fetch_city и при падении save_listings done_buckets == [], статус failed. Контроль «полный проход несёт унаследованное ∪ пройденное» — test_3369. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
248 lines
9.4 KiB
Python
248 lines
9.4 KiB
Python
"""Крайние случаи city-свипов kit'а, не покрытые после удаления legacy-пайплайна (#2406).
|
||
|
||
Уже покрыто: avito — таймаут якоря и дрейн после первого якоря (test_3319); дрейн ДО
|
||
первого якоря у yandex/cian/newbuilding (test_3333); отмена у domclick (test_3369).
|
||
Здесь остаток:
|
||
(а) таймаут якоря у cian и yandex; у domclick снятая/упавшая фаза не отмечает
|
||
корзины, строки которых не сохранены;
|
||
(б) SIGTERM-дрейн ПОСЛЕ первого якоря у cian и yandex;
|
||
(в) пользовательская отмена (runs.is_cancelled=True) после первого якоря у avito,
|
||
cian и yandex.
|
||
|
||
Проверяется строка прогона, а не вызовы: _FakeDb мержит counters всех записей по
|
||
порядку (как `counters || :counters` в runs.py) и запоминает финальный статус.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import os
|
||
|
||
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
|
||
from scraper_kit.orchestration import pipeline as pl
|
||
|
||
ANCHOR_A = (56.83, 60.60, "ekb-center")
|
||
ANCHOR_B = (56.79, 60.63, "ekb-south")
|
||
|
||
|
||
class _FakeDb:
|
||
def __init__(self) -> None:
|
||
self.counters: dict[str, Any] = {}
|
||
self.status = "running"
|
||
|
||
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
||
sql = str(stmt)
|
||
if params and "counters" in params:
|
||
self.counters.update(json.loads(params["counters"]))
|
||
for st in ("done", "failed", "banned"):
|
||
if f"SET status = '{st}'" in sql and self.status == "running":
|
||
self.status = st
|
||
return MagicMock()
|
||
|
||
def commit(self) -> None: ...
|
||
def rollback(self) -> None: ...
|
||
|
||
|
||
class _Scraper:
|
||
"""Двойник Cian/Yandex/AvitoScraper: помнит якоря, роняет заданный TimeoutError'ом."""
|
||
|
||
visited: list[float] = [] # noqa: RUF012
|
||
timeout_on: float | None = None
|
||
|
||
gate_fetch_attempts = state_extraction_attempts = 1
|
||
gate_fetch_failures = state_extraction_failures = 0
|
||
request_delay_sec = 0.0
|
||
|
||
def __init__(self, *_a: Any, **_kw: Any) -> None:
|
||
self._browser = None
|
||
self._cffi = None
|
||
|
||
async def __aenter__(self) -> _Scraper:
|
||
return self
|
||
|
||
async def __aexit__(self, *_e: Any) -> None:
|
||
return None
|
||
|
||
def _visit(self, lat: float) -> None:
|
||
_Scraper.visited.append(lat)
|
||
if lat == _Scraper.timeout_on:
|
||
raise TimeoutError
|
||
|
||
async def fetch_around_multi_room(self, lat: float, *_a: Any, **kw: Any) -> list[Any]:
|
||
self._visit(lat)
|
||
if kw.get("on_combo") is not None: # yandex: единица чекпоинта — combo
|
||
kw["on_combo"](f"combo@{lat}", [])
|
||
return []
|
||
return [types.SimpleNamespace(listing_segment="vtorichka", house_ext_id=None)]
|
||
|
||
async def fetch_around(self, lat: float, *_a: Any, **_kw: Any) -> list[Any]:
|
||
self._visit(lat)
|
||
return []
|
||
|
||
|
||
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,
|
||
)
|
||
|
||
|
||
async def _sweep(
|
||
kind: str,
|
||
*,
|
||
timeout_on: float | None = None,
|
||
drain_after_first: bool = False,
|
||
cancel_after_first: bool = False,
|
||
) -> _FakeDb:
|
||
_Scraper.visited = []
|
||
_Scraper.timeout_on = timeout_on
|
||
db = _FakeDb()
|
||
common: dict[str, Any] = {
|
||
"run_id": 2406,
|
||
"config": _config(),
|
||
"matcher": MagicMock(),
|
||
"anchors": [ANCHOR_A, ANCHOR_B],
|
||
"request_delay_sec": 0.0,
|
||
"shutdown_requested": lambda: drain_after_first and bool(_Scraper.visited),
|
||
}
|
||
with (
|
||
patch.object(pl, "CianScraper", _Scraper),
|
||
patch.object(pl, "YandexRealtyScraper", _Scraper),
|
||
patch.object(pl, "AvitoScraper", _Scraper),
|
||
patch.object(pl, "AsyncSession", _Scraper),
|
||
patch.object(pl, "save_listings", lambda *_a, **_kw: (1, 0)),
|
||
patch.object(
|
||
pl.runs, "is_cancelled", lambda *_a: cancel_after_first and bool(_Scraper.visited)
|
||
),
|
||
):
|
||
if kind == "cian":
|
||
await pl.run_cian_city_sweep(
|
||
db, # type: ignore[arg-type]
|
||
enrich_houses=False,
|
||
detail_top_n=0,
|
||
**common,
|
||
)
|
||
elif kind == "yandex":
|
||
enrichment = MagicMock()
|
||
enrichment.record_yandex_price_history.return_value = 0
|
||
await pl.run_yandex_city_sweep(
|
||
db, # type: ignore[arg-type]
|
||
enrichment=enrichment,
|
||
enrich_address=False,
|
||
**common,
|
||
)
|
||
else:
|
||
await pl.run_avito_city_sweep(
|
||
db, # type: ignore[arg-type]
|
||
enrichment=MagicMock(),
|
||
enrich_houses=False,
|
||
enrich_imv=False,
|
||
detail_top_n=0,
|
||
**common,
|
||
)
|
||
return db
|
||
|
||
|
||
def _done_key(kind: str, anchor: tuple[float, float, str]) -> str:
|
||
return f"combo@{anchor[0]}" if kind == "yandex" else anchor[2]
|
||
|
||
|
||
# ── (а) таймаут якоря: якорь не засчитан, следующий обработан ──────────────────
|
||
|
||
|
||
@pytest.mark.parametrize("kind", ["cian", "yandex"])
|
||
async def test_anchor_timeout_skips_anchor_and_continues(kind: str) -> None:
|
||
db = await _sweep(kind, timeout_on=ANCHOR_A[0])
|
||
|
||
assert _Scraper.visited == [ANCHOR_A[0], ANCHOR_B[0]]
|
||
assert db.counters["done_buckets"] == [_done_key(kind, ANCHOR_B)]
|
||
assert db.counters["errors_count"] >= 1
|
||
|
||
|
||
@pytest.mark.parametrize("broken", ["fetch_timeout", "save_failed"])
|
||
async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> None:
|
||
"""Корзина без сохранённых строк не попадает в чекпоинт.
|
||
|
||
Лоты Домклика копятся в памяти и пишутся ОДНИМ save_listings после всех корзин.
|
||
Фаза, снятая watchdog'ом (или упавшая на save), не сохранила ничего — даже из
|
||
корзины 'st', которую скрейпер успел пройти. Отметить её пройденной значило бы,
|
||
что следующий прогон пропустит её через skip_buckets навсегда (миграция 308).
|
||
"""
|
||
|
||
class _Dc(_Scraper):
|
||
blocked = False
|
||
geo_filtered = fetch_errors = bucket_start_index = 0
|
||
buckets_completed, buckets_total = 1, 6
|
||
completed_buckets = ["st"] # noqa: RUF012
|
||
|
||
async def fetch_city(self, **_kw: Any) -> list[Any]:
|
||
if broken == "fetch_timeout":
|
||
raise TimeoutError
|
||
return [object()]
|
||
|
||
saved: list[int] = []
|
||
|
||
def _save(_db: Any, lots: list[Any], **_kw: Any) -> tuple[int, int]:
|
||
if broken == "save_failed":
|
||
raise RuntimeError("save_listings упал")
|
||
saved.append(len(lots))
|
||
return len(lots), 0
|
||
|
||
db = _FakeDb()
|
||
with (
|
||
patch.object(pl, "DomClickScraper", _Dc),
|
||
patch.object(pl, "save_listings", _save),
|
||
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
||
):
|
||
counters = await pl.run_domclick_city_sweep(
|
||
db, # type: ignore[arg-type]
|
||
run_id=2406,
|
||
config=types.SimpleNamespace(browser_http_endpoint="http://x:9000"),
|
||
matcher=MagicMock(),
|
||
pages=1,
|
||
request_delay_sec=0.0,
|
||
)
|
||
|
||
assert saved == [], "тест не о том: строки сохранились"
|
||
assert counters.errors_count >= 1
|
||
assert db.counters["done_buckets"] == [], (
|
||
f"корзины без единой сохранённой строки в чекпоинте: {db.counters['done_buckets']}"
|
||
)
|
||
assert db.status == "failed"
|
||
|
||
|
||
# ── (б) дрейн после первого якоря: interrupted=1, в чекпоинте ровно первый ──────
|
||
|
||
|
||
@pytest.mark.parametrize("kind", ["cian", "yandex"])
|
||
async def test_drain_after_first_anchor_is_partial_and_interrupted(kind: str) -> None:
|
||
db = await _sweep(kind, drain_after_first=True)
|
||
|
||
assert _Scraper.visited == [ANCHOR_A[0]]
|
||
assert db.counters["interrupted"] == 1
|
||
assert db.counters["done_buckets"] == [_done_key(kind, ANCHOR_A)]
|
||
assert db.status == "done"
|
||
|
||
|
||
# ── (в) пользовательская отмена после первого якоря ────────────────────────────
|
||
|
||
|
||
@pytest.mark.parametrize("kind", ["avito", "cian", "yandex"])
|
||
async def test_user_cancel_after_first_anchor_stops_without_finalizing(kind: str) -> None:
|
||
db = await _sweep(kind, cancel_after_first=True)
|
||
|
||
assert _Scraper.visited == [ANCHOR_A[0]]
|
||
assert db.counters["done_buckets"] == [_done_key(kind, ANCHOR_A)]
|
||
assert "interrupted" not in db.counters
|
||
# Строку финализирует сам mark_cancelled (UI); свип только останавливается.
|
||
assert db.status == "running"
|