fix(scrapers): cian full_load anti-zombie + avito changeip-resilience & честный статус (#1949 #1950) #1952

Merged
lekss361 merged 2 commits from fix/scraper-resilience-cian-zombie-avito-changeip into main 2026-06-27 08:38:11 +00:00
4 changed files with 711 additions and 28 deletions

View file

@ -359,6 +359,17 @@ class Settings(BaseSettings):
# Settle-sleep после changeip-вызова: мобильный модем поднимает новый IP. # Settle-sleep после changeip-вызова: мобильный модем поднимает новый IP.
# ~9с по умолчанию (эмпирика mobileproxy.space). ENV: AVITO_PROXY_ROTATE_SETTLE_S. # ~9с по умолчанию (эмпирика mobileproxy.space). ENV: AVITO_PROXY_ROTATE_SETTLE_S.
avito_proxy_rotate_settle_s: float = 9.0 avito_proxy_rotate_settle_s: float = 9.0
# #1950: retry-параметры changeip-GET (_rotate_proxy_ip). Вместо одношотного 30s-timeout
# делаем proxy_rotate_attempts попыток по proxy_rotate_attempt_timeout_s каждая.
# Короткий timeout (8s) означает, что зависший changeip не блокирует весь run на 30s.
# ENV: PROXY_ROTATE_ATTEMPT_TIMEOUT_S / PROXY_ROTATE_ATTEMPTS.
proxy_rotate_attempt_timeout_s: float = 8.0
proxy_rotate_attempts: int = 3
# #1950: если SERP уже сохранил лоты (ins+upd > 0) и упали только detail/houses,
# ставим 'done' а не 'banned' — partial intake сохранён, 'banned' лишний.
# False = старое поведение. ENV: AVITO_SERP_OK_NOT_BANNED.
avito_serp_ok_not_banned: bool = True
# ── Cian dedicated mobile proxy (separate egress from Avito) ────────────── # ── Cian dedicated mobile proxy (separate egress from Avito) ──────────────
# Cian и Avito делят один мобильный IP при общем scraper_proxy_url → конкуренция # Cian и Avito делят один мобильный IP при общем scraper_proxy_url → конкуренция
@ -373,6 +384,13 @@ class Settings(BaseSettings):
# ENV: CIAN_PROXY_MAX_ROTATIONS. # ENV: CIAN_PROXY_MAX_ROTATIONS.
cian_proxy_max_rotations: int = 4 cian_proxy_max_rotations: int = 4
# #1949: per-fetch hard-cancel для browser-fetch'ей внутри bucket-gather Cian full_load.
# asyncio.wait_for(fetch, timeout=X) отменяет зависший fetch → TimeoutError →
# _fetch_page_html ловит как Exception → None → _one_page → [] → gather завершается.
# 0 = отключить (старое поведение). Дефолт 90.0s щедрее nормального fetch ~12s.
# ENV: CIAN_FULL_LOAD_PER_FETCH_TIMEOUT_S.
cian_full_load_per_fetch_timeout_s: float = 90.0
@property @property
def cian_proxy_url(self) -> str | None: def cian_proxy_url(self) -> str | None:
"""Прокси для Cian-скраперов. CIAN_PROXY_URL > scraper_proxy_url (fallback).""" """Прокси для Cian-скраперов. CIAN_PROXY_URL > scraper_proxy_url (fallback)."""

View file

@ -33,6 +33,7 @@ import random
from contextlib import AsyncExitStack from contextlib import AsyncExitStack
from dataclasses import dataclass, field, fields from dataclasses import dataclass, field, fields
from datetime import date, timedelta from datetime import date, timedelta
from typing import Any
from urllib.parse import urlparse from urllib.parse import urlparse
from curl_cffi.requests import AsyncSession from curl_cffi.requests import AsyncSession
@ -81,35 +82,55 @@ async def _rotate_proxy_ip(*, reason: str, rotations_done: int, source: str = "a
if not rotate_url: if not rotate_url:
return False return False
sep = "&" if "?" in rotate_url else "?" sep = "&" if "?" in rotate_url else "?"
try: # #1950: retry changeip-GET короткими попытками вместо одношотного 30s-timeout.
async with AsyncSession(timeout=30) as rot: # Если changeip-сервер завис на 30s → весь прогон ждёт зря + ротация считается
resp = await rot.get(f"{rotate_url}{sep}format=json") # успешной (False возвращался). Теперь: proxy_rotate_attempts попыток по
new_ip: str = "" # proxy_rotate_attempt_timeout_s каждая; возвращаем True при первом успехе.
_attempts = settings.proxy_rotate_attempts
_attempt_timeout = settings.proxy_rotate_attempt_timeout_s
for attempt in range(_attempts):
try: try:
data = resp.json() async with AsyncSession(timeout=_attempt_timeout) as rot:
new_ip = str(data.get("new_ip", "")) resp = await rot.get(f"{rotate_url}{sep}format=json")
new_ip: str = ""
try:
data = resp.json()
new_ip = str(data.get("new_ip", ""))
except Exception:
pass
await asyncio.sleep(settings.avito_proxy_rotate_settle_s)
logger.info(
"pipeline: IP rotated via changeip — source=%s reason=%s "
"rotation=#%d attempt=%d/%d remaining=%d new_ip=%s",
source,
reason,
rotations_done + 1,
attempt + 1,
_attempts,
max_rot - rotations_done - 1,
new_ip or "unknown",
)
return True
except Exception: except Exception:
pass logger.warning(
await asyncio.sleep(settings.avito_proxy_rotate_settle_s) "pipeline: IP rotation attempt %d/%d failed " "(source=%s reason=%s rotation=#%d)",
logger.info( attempt + 1,
"pipeline: IP rotated via changeip — source=%s reason=%s " _attempts,
"rotation=#%d remaining=%d new_ip=%s", source,
source, reason,
reason, rotations_done + 1,
rotations_done + 1, exc_info=True,
max_rot - rotations_done - 1, )
new_ip or "unknown", if attempt < _attempts - 1:
) await asyncio.sleep(1.0)
return True logger.error(
except Exception: "pipeline: all %d changeip attempts failed " "(source=%s reason=%s rotation=#%d)",
logger.warning( _attempts,
"pipeline: IP rotation failed (source=%s reason=%s rotation=#%d)", source,
source, reason,
reason, rotations_done + 1,
rotations_done + 1, )
exc_info=True, return False
)
return False
# Per-anchor watchdog timeout (секунды). Если обработка одного anchor'а занимает # Per-anchor watchdog timeout (секунды). Если обработка одного anchor'а занимает
@ -1187,7 +1208,30 @@ async def run_avito_city_sweep(
counters.errors_count += 1 counters.errors_count += 1
counters.anchors_done = idx counters.anchors_done = idx
scrape_runs.update_heartbeat(db, run_id, counters.to_dict()) scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
scrape_runs.mark_banned(db, run_id, str(e), counters.to_dict()) # #1950: если SERP уже собрал лоты и заблокировало только detail/houses,
# ставим 'done' (не 'banned') — partial intake сохранён.
# За флагом avito_serp_ok_not_banned (default True).
_serp_ok = (counters.lots_inserted + counters.lots_updated) > 0
if settings.avito_serp_ok_not_banned and _serp_ok:
_note = (
f"detail enrichment aborted ({e}) но SERP intake завершён: "
f"{counters.lots_inserted} ins {counters.lots_updated} upd"
)
logger.warning(
"city-sweep run_id=%d: SERP OK (ins=%d upd=%d), "
"detail blocked — DONE not banned. %s",
run_id,
counters.lots_inserted,
counters.lots_updated,
_note,
)
scrape_runs.mark_done(
db,
run_id,
{**counters.to_dict(), "enrichment_abort_note": _note}, # type: ignore[arg-type]
)
else:
scrape_runs.mark_banned(db, run_id, str(e), counters.to_dict())
return counters return counters
except Exception: except Exception:
logger.exception("city-sweep run_id=%d: anchor %s failed", run_id, name) logger.exception("city-sweep run_id=%d: anchor %s failed", run_id, name)
@ -2492,10 +2536,61 @@ async def run_cian_full_load(
"""Heartbeat per room-bucket (на случай пустых бакетов без on_bucket вызовов).""" """Heartbeat per room-bucket (на случай пустых бакетов без on_bucket вызовов)."""
scrape_runs.update_heartbeat(db, run_id, counters.to_dict()) scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
# #1949: фоновый heartbeat — обновляет heartbeat_at каждые 60s независимо от on_bucket.
# Без него зависший asyncio.gather внутри bucket-loop замораживает heartbeat на часы
# → reap_zombies помечает живой run zombie. Отменяется в finally при любом завершении.
async def _background_heartbeat() -> None:
"""Периодический heartbeat каждые 60s, не завязанный на прогресс бакетов."""
while True:
await asyncio.sleep(60.0)
try:
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
except Exception:
pass
_hb_task: asyncio.Task[None] | None = None
try: try:
_hb_task = asyncio.create_task(_background_heartbeat())
async with CianScraper() as scraper: async with CianScraper() as scraper:
scraper.request_delay_sec = request_delay_sec scraper.request_delay_sec = request_delay_sec
# #1949: per-fetch hard-cancel для browser-fetch'ей внутри bucket-gather.
# BrowserFetcher.fetch имеет httpx-timeout 120s, но browser-сервис иногда
# виснет так что httpx не получает ответ → gather морозит → heartbeat стоп.
# Оборачиваем _browser.fetch в asyncio.wait_for(timeout=X) через monkey-patch:
# TimeoutError → _fetch_page_html ловит как Exception → None → _one_page → []
# → gather завершается нормально → on_bucket вызывается → heartbeat обновляется.
# За флагом cian_full_load_per_fetch_timeout_s (0 = отключить).
_fetch_timeout = settings.cian_full_load_per_fetch_timeout_s
if _fetch_timeout > 0 and scraper._browser is not None:
_orig_fetch = scraper._browser.fetch
async def _timed_fetch(
url: str,
_f: Any = _orig_fetch,
_t: float = _fetch_timeout,
_rid: int = run_id,
) -> str:
try:
return await asyncio.wait_for(_f(url), timeout=_t)
except TimeoutError:
logger.warning(
"cian-full-load run_id=%d: browser-fetch timed out "
"after %.0fs — url=%s",
_rid,
_t,
url,
)
raise
scraper._browser.fetch = _timed_fetch # type: ignore[method-assign]
logger.info(
"cian-full-load run_id=%d: per-fetch timeout %.0fs enabled",
run_id,
_fetch_timeout,
)
await scraper.fetch_all_secondary( await scraper.fetch_all_secondary(
price_cap_per_bucket=price_cap_per_bucket, price_cap_per_bucket=price_cap_per_bucket,
concurrency=concurrency, concurrency=concurrency,
@ -2609,6 +2704,15 @@ async def run_cian_full_load(
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict()) scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
raise raise
finally:
# #1949: отменяем фоновый heartbeat при любом завершении (return / raise / cancel).
if _hb_task is not None:
_hb_task.cancel()
try:
await _hb_task
except asyncio.CancelledError:
pass
# ── Yandex exhaustive full load (region-wide, без anchor'ов) ────────────────── # ── Yandex exhaustive full load (region-wide, без anchor'ов) ──────────────────

View file

@ -200,6 +200,7 @@ async def test_rotate_proxy_ip_returns_false_on_network_error() -> None:
with ( with (
patch("app.services.scrape_pipeline.settings") as s, patch("app.services.scrape_pipeline.settings") as s,
patch("app.services.scrape_pipeline.AsyncSession", mock_session_cls), patch("app.services.scrape_pipeline.AsyncSession", mock_session_cls),
patch("app.services.scrape_pipeline.asyncio.sleep", AsyncMock()),
): ):
s.cian_proxy_rotate_url = "http://cian-rotate.test/" s.cian_proxy_rotate_url = "http://cian-rotate.test/"
s.avito_proxy_rotate_url = None s.avito_proxy_rotate_url = None
@ -208,6 +209,8 @@ async def test_rotate_proxy_ip_returns_false_on_network_error() -> None:
s.avito_proxy_max_rotations = 4 s.avito_proxy_max_rotations = 4
s.yandex_proxy_max_rotations = 4 s.yandex_proxy_max_rotations = 4
s.avito_proxy_rotate_settle_s = 0.0 s.avito_proxy_rotate_settle_s = 0.0
s.proxy_rotate_attempts = 1 # одна попытка, без retry-sleep (#1950)
s.proxy_rotate_attempt_timeout_s = 8.0
result = await _rotate_proxy_ip(reason="ban", rotations_done=0, source="cian") result = await _rotate_proxy_ip(reason="ban", rotations_done=0, source="cian")

View file

@ -0,0 +1,558 @@
"""Тесты устойчивости скрейперов: avito changeip-retry (#1950) + cian anti-zombie (#1949).
Fix A (#1950): _rotate_proxy_ip теперь делает до proxy_rotate_attempts попыток
по proxy_rotate_attempt_timeout_s каждая.
Fix B (#1950): run_avito_city_sweep при AvitoBlockedError/RateLimited с
lots_inserted+lots_updated>0 вызывает mark_done (не mark_banned)
когда avito_serp_ok_not_banned=True.
Fix C (#1949): run_cian_full_load применяет monkey-patch _browser.fetch с
asyncio.wait_for(timeout=cian_full_load_per_fetch_timeout_s).
Fix D (#1949): run_cian_full_load создаёт фоновый heartbeat-task и отменяет
его в finally при любом завершении.
"""
from __future__ import annotations
import asyncio
import os
import time
from typing import Any
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
from app.services.scrapers.avito_exceptions import AvitoBlockedError
from app.services.scrapers.base import ScrapedLot
# ──────────────────────────────────────────────────────────────────────────────
# Fix A: _rotate_proxy_ip retry loop (#1950)
# ──────────────────────────────────────────────────────────────────────────────
def _make_session_mock(get_side_effect: Any) -> AsyncMock:
"""Вспомогательный: AsyncMock-сессия с заданным поведением .get()."""
mock_session = AsyncMock()
mock_session.__aenter__ = AsyncMock(return_value=mock_session)
mock_session.__aexit__ = AsyncMock(return_value=None)
mock_session.get = AsyncMock(side_effect=get_side_effect)
return mock_session
@pytest.mark.asyncio
async def test_changeip_retry_success_on_second_attempt() -> None:
"""#1950-A: первая попытка таймаут, вторая успех → возвращает True; ровно 2 GET-вызова."""
from app.services.scrape_pipeline import _rotate_proxy_ip
get_call_count = 0
async def fake_get(url: str) -> MagicMock:
nonlocal get_call_count
get_call_count += 1
if get_call_count == 1:
raise TimeoutError("connect timeout")
resp = MagicMock()
resp.json.return_value = {"new_ip": "1.2.3.4"}
return resp
mock_session = _make_session_mock(fake_get)
with (
patch("app.services.scrape_pipeline.AsyncSession", return_value=mock_session),
patch("app.services.scrape_pipeline.asyncio.sleep", AsyncMock()),
patch("app.services.scrape_pipeline.settings") as s,
):
s.proxy_rotate_attempts = 3
s.proxy_rotate_attempt_timeout_s = 8.0
s.avito_proxy_rotate_settle_s = 0.0
s.avito_proxy_rotate_url = "http://changeip.example.com/rotate?token=X"
s.avito_proxy_max_rotations = 3
result = await _rotate_proxy_ip(
source="avito",
reason="test-detail-block",
rotations_done=0,
)
assert result is True
assert get_call_count == 2, f"expected 2 GET calls, got {get_call_count}"
@pytest.mark.asyncio
async def test_changeip_retry_all_fail_returns_false() -> None:
"""#1950-A: все 3 попытки падают с TimeoutError → False; ровно 3 GET-вызова."""
from app.services.scrape_pipeline import _rotate_proxy_ip
get_call_count = 0
async def failing_get(url: str) -> None:
nonlocal get_call_count
get_call_count += 1
raise TimeoutError("connect timeout")
mock_session = _make_session_mock(failing_get)
with (
patch("app.services.scrape_pipeline.AsyncSession", return_value=mock_session),
patch("app.services.scrape_pipeline.asyncio.sleep", AsyncMock()),
patch("app.services.scrape_pipeline.settings") as s,
):
s.proxy_rotate_attempts = 3
s.proxy_rotate_attempt_timeout_s = 8.0
s.avito_proxy_rotate_settle_s = 0.0
s.avito_proxy_rotate_url = "http://changeip.example.com/rotate?token=X"
s.avito_proxy_max_rotations = 3
result = await _rotate_proxy_ip(
source="avito",
reason="test-detail-block",
rotations_done=0,
)
assert result is False
assert (
get_call_count == 3
), f"expected exactly 3 attempts (proxy_rotate_attempts=3), got {get_call_count}"
@pytest.mark.asyncio
async def test_changeip_retry_bounded_by_setting() -> None:
"""#1950-A: proxy_rotate_attempts=2 → ровно 2 попытки, не больше."""
from app.services.scrape_pipeline import _rotate_proxy_ip
call_count = 0
async def failing_get(url: str) -> None:
nonlocal call_count
call_count += 1
raise TimeoutError("timeout")
mock_session = _make_session_mock(failing_get)
with (
patch("app.services.scrape_pipeline.AsyncSession", return_value=mock_session),
patch("app.services.scrape_pipeline.asyncio.sleep", AsyncMock()),
patch("app.services.scrape_pipeline.settings") as s,
):
s.proxy_rotate_attempts = 2
s.proxy_rotate_attempt_timeout_s = 8.0
s.avito_proxy_rotate_settle_s = 0.0
s.cian_proxy_rotate_url = "http://changeip.example.com/rotate"
s.avito_proxy_rotate_url = None
s.cian_proxy_max_rotations = 3
result = await _rotate_proxy_ip(
source="cian",
reason="serp-block",
rotations_done=1,
)
assert result is False
assert call_count == 2, f"proxy_rotate_attempts=2 should make exactly 2 calls, got {call_count}"
# ──────────────────────────────────────────────────────────────────────────────
# Fix B: честный статус run_avito_city_sweep (#1950)
# ──────────────────────────────────────────────────────────────────────────────
def _make_city_sweep_mocks() -> tuple[MagicMock, MagicMock, MagicMock]:
"""Возвращает (mock_db, mock_runs, mock_scraper) для тестов run_avito_city_sweep."""
mock_db = MagicMock()
# priority_rows для detail-фазы: 3 строки с source_url
class _DetailRow:
def __getitem__(self, key: str) -> str:
return "https://www.avito.ru/ekaterinburg/kvartiry/test-1"
mock_db.execute.return_value.mappings.return_value.all.return_value = [
_DetailRow(),
_DetailRow(),
_DetailRow(),
]
mock_runs = MagicMock()
mock_runs.is_cancelled.return_value = False
mock_runs.mark_done = MagicMock()
mock_runs.mark_banned = MagicMock()
mock_runs.mark_failed = MagicMock()
mock_runs.update_heartbeat = MagicMock()
mock_scraper = MagicMock()
# AvitoScraper НЕ используется как async CM в _avito_anchor_phases (просто instantiate)
mock_scraper.fetch_around = AsyncMock(
return_value=[
ScrapedLot(
source="avito",
source_url=f"https://www.avito.ru/e/k/{i}",
source_id=str(i),
price_rub=5_000_000,
address="ЕКБ test",
listing_segment="vtorichka",
)
for i in range(3)
]
)
return mock_db, mock_runs, mock_scraper
@pytest.mark.asyncio
async def test_avito_city_sweep_serp_ok_marks_done_not_banned() -> None:
"""#1950-B: SERP собрал лоты (ins=3), detail заблокировало → mark_done не mark_banned.
Условие: avito_serp_ok_not_banned=True (default), lots_inserted+lots_updated > 0.
"""
from app.services.scrape_pipeline import run_avito_city_sweep
mock_db, mock_runs, mock_scraper = _make_city_sweep_mocks()
# Shared AsyncSession mock (browser_mode=False)
mock_session = AsyncMock()
with (
patch("app.services.scrape_pipeline.AvitoScraper", return_value=mock_scraper),
patch("app.services.scrape_pipeline.AsyncSession", return_value=mock_session),
patch("app.services.scrape_pipeline.scrape_runs", mock_runs),
patch("app.services.scrape_pipeline.save_listings", return_value=(3, 0)),
patch(
"app.services.scrape_pipeline.fetch_detail",
AsyncMock(side_effect=AvitoBlockedError("HTTP 403 banned")),
),
patch("app.services.scrape_pipeline.asyncio.sleep", AsyncMock()),
patch("app.services.scrape_pipeline.settings") as s,
):
s.scraper_fetch_mode = "curl_cffi"
s.scraper_proxy_url = None
s.avito_proxy_max_rotations = 0 # no rotation budget
s.avito_serp_ok_not_banned = True # ← наш флаг
s.avito_proxy_rotate_settle_s = 0.0
s.proxy_rotate_attempts = 3
s.proxy_rotate_attempt_timeout_s = 8.0
await run_avito_city_sweep(
mock_db,
run_id=1,
anchors=[(56.84, 60.61, "test-anchor")],
detail_top_n=3,
enrich_houses=False,
enrich_imv=False,
)
assert (
mock_runs.mark_done.called
), "mark_done должен быть вызван: SERP OK (ins=3), detail blocked"
assert (
not mock_runs.mark_banned.called
), "mark_banned НЕ должен быть вызван когда SERP уже сохранил лоты"
# enrichment_abort_note передан в counters
call_args = mock_runs.mark_done.call_args
# mark_done(db, run_id, counters_dict)
passed_counters: dict = call_args[0][2] if call_args[0] else {}
assert (
"enrichment_abort_note" in passed_counters
), "enrichment_abort_note должен присутствовать в counters переданных mark_done"
@pytest.mark.asyncio
async def test_avito_city_sweep_serp_zero_lots_marks_banned() -> None:
"""#1950-B: SERP вернул 0 лотов, SERP-блок → mark_banned (не mark_done).
Условие: lots_inserted+lots_updated == 0.
"""
from app.services.scrape_pipeline import run_avito_city_sweep
mock_db, mock_runs, _ = _make_city_sweep_mocks()
# SERP сам поднимает AvitoBlockedError → anchor_lots не сохраняется
mock_scraper_serp_block = MagicMock()
mock_scraper_serp_block.fetch_around = AsyncMock(side_effect=AvitoBlockedError("SERP HTTP 403"))
mock_session = AsyncMock()
with (
patch("app.services.scrape_pipeline.AvitoScraper", return_value=mock_scraper_serp_block),
patch("app.services.scrape_pipeline.AsyncSession", return_value=mock_session),
patch("app.services.scrape_pipeline.scrape_runs", mock_runs),
patch("app.services.scrape_pipeline.asyncio.sleep", AsyncMock()),
patch("app.services.scrape_pipeline.settings") as s,
):
s.scraper_fetch_mode = "curl_cffi"
s.scraper_proxy_url = None
s.avito_proxy_max_rotations = 0
s.avito_serp_ok_not_banned = True # включён флаг, но лотов нет
s.avito_proxy_rotate_settle_s = 0.0
s.proxy_rotate_attempts = 3
s.proxy_rotate_attempt_timeout_s = 8.0
await run_avito_city_sweep(
mock_db,
run_id=2,
anchors=[(56.84, 60.61, "test-anchor")],
detail_top_n=3,
enrich_houses=False,
enrich_imv=False,
)
assert (
mock_runs.mark_banned.called
), "mark_banned должен быть вызван: SERP сам заблокирован, lots=0"
assert not mock_runs.mark_done.called, "mark_done НЕ должен быть вызван при lots_inserted == 0"
@pytest.mark.asyncio
async def test_avito_city_sweep_flag_off_still_marks_banned() -> None:
"""#1950-B: avito_serp_ok_not_banned=False → mark_banned даже при lots>0 (backward-compat)."""
from app.services.scrape_pipeline import run_avito_city_sweep
mock_db, mock_runs, mock_scraper = _make_city_sweep_mocks()
mock_session = AsyncMock()
with (
patch("app.services.scrape_pipeline.AvitoScraper", return_value=mock_scraper),
patch("app.services.scrape_pipeline.AsyncSession", return_value=mock_session),
patch("app.services.scrape_pipeline.scrape_runs", mock_runs),
patch("app.services.scrape_pipeline.save_listings", return_value=(3, 0)),
patch(
"app.services.scrape_pipeline.fetch_detail",
AsyncMock(side_effect=AvitoBlockedError("HTTP 403")),
),
patch("app.services.scrape_pipeline.asyncio.sleep", AsyncMock()),
patch("app.services.scrape_pipeline.settings") as s,
):
s.scraper_fetch_mode = "curl_cffi"
s.scraper_proxy_url = None
s.avito_proxy_max_rotations = 0
s.avito_serp_ok_not_banned = False # ← флаг выключен
s.avito_proxy_rotate_settle_s = 0.0
s.proxy_rotate_attempts = 3
s.proxy_rotate_attempt_timeout_s = 8.0
await run_avito_city_sweep(
mock_db,
run_id=3,
anchors=[(56.84, 60.61, "test-anchor")],
detail_top_n=3,
enrich_houses=False,
enrich_imv=False,
)
assert (
mock_runs.mark_banned.called
), "avito_serp_ok_not_banned=False → mark_banned (backward-compat)"
assert not mock_runs.mark_done.called
# ──────────────────────────────────────────────────────────────────────────────
# Fix C: cian per-fetch wait_for (#1949)
# ──────────────────────────────────────────────────────────────────────────────
@pytest.mark.asyncio
async def test_cian_full_load_hanging_fetch_cancelled_by_timeout() -> None:
"""#1949-C: browser-fetch зависший навсегда отменяется per-fetch timeout.
Monkey-patch заменяет _browser.fetch на _timed_fetch(wait_for timeout=0.05s).
Зависший fetch отменяется через ~50ms TimeoutError run завершается.
"""
from app.services.scrape_pipeline import run_cian_full_load
fetch_was_called = False
async def hanging_fetch(url: str) -> str:
nonlocal fetch_was_called
fetch_was_called = True
await asyncio.Event().wait() # висит вечно
return ""
mock_browser = MagicMock()
mock_browser.fetch = hanging_fetch
mock_scraper = MagicMock()
mock_scraper.__aenter__ = AsyncMock(return_value=mock_scraper)
mock_scraper.__aexit__ = AsyncMock(return_value=None)
mock_scraper._browser = mock_browser
mock_scraper.request_delay_sec = 0.0
# fetch_all_secondary симулирует _fetch_page_html: вызывает fetch и ловит Exception
async def fake_fetch_all_secondary(**kwargs: object) -> None:
try:
# После monkey-patch _browser.fetch = _timed_fetch (с timeout=0.05s)
await mock_scraper._browser.fetch("https://ekb.cian.ru/sale/flat/1/")
except Exception:
pass # TimeoutError → аналог _fetch_page_html → returns None
mock_scraper.fetch_all_secondary = fake_fetch_all_secondary
mock_db = MagicMock()
mock_runs = MagicMock()
mock_runs.is_cancelled.return_value = False
mock_runs.mark_done = MagicMock()
mock_runs.mark_failed = MagicMock()
mock_runs.update_heartbeat = MagicMock()
t0 = time.monotonic()
with (
patch("app.services.scrapers.cian.CianScraper", return_value=mock_scraper),
patch("app.services.scrape_pipeline.scrape_runs", mock_runs),
patch("app.services.scrape_pipeline.settings") as s,
):
s.cian_full_load_per_fetch_timeout_s = 0.05 # 50ms — быстрый timeout для теста
s.cian_full_load_detail_top_n = 0
await run_cian_full_load(mock_db, run_id=1)
elapsed = time.monotonic() - t0
assert fetch_was_called, "hanging_fetch должен быть вызван через monkey-patched _browser.fetch"
assert (
elapsed < 3.0
), f"run должен завершиться быстро (fetch timeout 50ms), занял {elapsed:.2f}s"
assert mock_runs.mark_done.called, "run_cian_full_load должен завершиться mark_done"
@pytest.mark.asyncio
async def test_cian_full_load_per_fetch_timeout_disabled_when_zero() -> None:
"""#1949-C: cian_full_load_per_fetch_timeout_s=0 → monkey-patch НЕ применяется."""
from app.services.scrape_pipeline import run_cian_full_load
original_fetch = AsyncMock(return_value="<html/>")
mock_browser = MagicMock()
mock_browser.fetch = original_fetch
mock_scraper = MagicMock()
mock_scraper.__aenter__ = AsyncMock(return_value=mock_scraper)
mock_scraper.__aexit__ = AsyncMock(return_value=None)
mock_scraper._browser = mock_browser
mock_scraper.request_delay_sec = 0.0
mock_scraper.fetch_all_secondary = AsyncMock(return_value=None)
mock_db = MagicMock()
mock_runs = MagicMock()
mock_runs.is_cancelled.return_value = False
mock_runs.mark_done = MagicMock()
mock_runs.mark_failed = MagicMock()
mock_runs.update_heartbeat = MagicMock()
with (
patch("app.services.scrapers.cian.CianScraper", return_value=mock_scraper),
patch("app.services.scrape_pipeline.scrape_runs", mock_runs),
patch("app.services.scrape_pipeline.settings") as s,
):
s.cian_full_load_per_fetch_timeout_s = 0.0 # отключено
s.cian_full_load_detail_top_n = 0
await run_cian_full_load(mock_db, run_id=2)
# Когда timeout=0 (отключён), monkey-patch не применяется → fetch остаётся original
assert (
mock_browser.fetch is original_fetch
), "cian_full_load_per_fetch_timeout_s=0 → fetch НЕ должен быть заменён monkey-patch"
# ──────────────────────────────────────────────────────────────────────────────
# Fix D: cian background heartbeat created and cancelled (#1949)
# ──────────────────────────────────────────────────────────────────────────────
@pytest.mark.asyncio
async def test_cian_full_load_background_heartbeat_created_and_cancelled() -> None:
"""#1949-D: background heartbeat task создаётся и отменяется в finally."""
from app.services.scrape_pipeline import run_cian_full_load
mock_scraper = MagicMock()
mock_scraper.__aenter__ = AsyncMock(return_value=mock_scraper)
mock_scraper.__aexit__ = AsyncMock(return_value=None)
mock_scraper._browser = None # per-fetch timeout не применяется
mock_scraper.fetch_all_secondary = AsyncMock(return_value=None)
mock_db = MagicMock()
mock_runs = MagicMock()
mock_runs.is_cancelled.return_value = False
mock_runs.mark_done = MagicMock()
mock_runs.mark_failed = MagicMock()
mock_runs.update_heartbeat = MagicMock()
created_tasks: list[asyncio.Task[None]] = []
_real_create_task = asyncio.create_task
def _tracking_create_task(coro: Any, *, name: str | None = None) -> asyncio.Task[None]:
task = _real_create_task(coro, name=name)
created_tasks.append(task)
return task # type: ignore[return-value]
with (
patch("app.services.scrapers.cian.CianScraper", return_value=mock_scraper),
patch("app.services.scrape_pipeline.scrape_runs", mock_runs),
patch(
"app.services.scrape_pipeline.asyncio.create_task",
side_effect=_tracking_create_task,
),
patch("app.services.scrape_pipeline.settings") as s,
):
s.cian_full_load_per_fetch_timeout_s = 0.0
s.cian_full_load_detail_top_n = 0
await run_cian_full_load(mock_db, run_id=3)
assert (
len(created_tasks) >= 1
), "asyncio.create_task должен быть вызван хотя бы раз (background heartbeat)"
for task in created_tasks:
assert (
task.cancelled() or task.done()
), f"Heartbeat task должен быть отменён в finally: {task}"
@pytest.mark.asyncio
async def test_cian_full_load_heartbeat_cancelled_on_exception() -> None:
"""#1949-D: heartbeat task отменяется в finally даже при исключении в run."""
from app.services.scrape_pipeline import run_cian_full_load
mock_scraper = MagicMock()
mock_scraper.__aenter__ = AsyncMock(return_value=mock_scraper)
mock_scraper.__aexit__ = AsyncMock(return_value=None)
mock_scraper._browser = None
mock_scraper.fetch_all_secondary = AsyncMock(side_effect=RuntimeError("fatal-test-error"))
mock_db = MagicMock()
mock_runs = MagicMock()
mock_runs.is_cancelled.return_value = False
mock_runs.mark_done = MagicMock()
mock_runs.mark_failed = MagicMock()
mock_runs.update_heartbeat = MagicMock()
created_tasks: list[asyncio.Task[None]] = []
_real_create_task = asyncio.create_task
def _tracking_create_task(coro: Any, *, name: str | None = None) -> asyncio.Task[None]:
task = _real_create_task(coro, name=name)
created_tasks.append(task)
return task # type: ignore[return-value]
with (
patch("app.services.scrapers.cian.CianScraper", return_value=mock_scraper),
patch("app.services.scrape_pipeline.scrape_runs", mock_runs),
patch(
"app.services.scrape_pipeline.asyncio.create_task",
side_effect=_tracking_create_task,
),
patch("app.services.scrape_pipeline.settings") as s,
):
s.cian_full_load_per_fetch_timeout_s = 0.0
s.cian_full_load_detail_top_n = 0
with pytest.raises(RuntimeError, match="fatal-test-error"):
await run_cian_full_load(mock_db, run_id=4)
assert len(created_tasks) >= 1, "heartbeat task должен быть создан"
for task in created_tasks:
assert (
task.cancelled() or task.done()
), f"Heartbeat task должен быть отменён в finally даже при RuntimeError: {task}"
assert mock_runs.mark_failed.called, "mark_failed должен быть вызван при RuntimeError"