fix(tradein/avito-detail): живой прогрев сессии + батч-на-сессию — снимает 403 на curl-backfill (#1551)
All checks were successful
Deploy Trade-In / changes (push) Successful in 12s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 1m27s
Deploy Trade-In / build-backend (push) Successful in 51s
Deploy Trade-In / deploy (push) Successful in 1m37s

This commit is contained in:
lekss361 2026-06-28 22:34:38 +00:00
parent cbe56cab0e
commit 69845ddcee
4 changed files with 431 additions and 26 deletions

View file

@ -22,11 +22,12 @@ from __future__ import annotations
import asyncio import asyncio
import json import json
import logging import logging
import random
import re import re
from dataclasses import dataclass, field from dataclasses import dataclass, field
from datetime import date from datetime import date
from typing import TYPE_CHECKING, Any from typing import TYPE_CHECKING, Any
from urllib.parse import urljoin, urlparse from urllib.parse import quote, urljoin, urlparse
from curl_cffi.requests import AsyncSession from curl_cffi.requests import AsyncSession
from selectolax.parser import HTMLParser, Node from selectolax.parser import HTMLParser, Node
@ -49,6 +50,18 @@ logger = logging.getLogger(__name__)
AVITO_BASE = "https://www.avito.ru" AVITO_BASE = "https://www.avito.ru"
# Живой прогрев сессии (#1551): GET Yandex SERP (органический referer-источник) -> GET
# Avito EKB search-страница (Referer=yandex) сеет антибот-куки, после чего прогретая
# curl_cffi-сессия держит батч detail-GET'ов без 403 (доказано пробами на проде: 22/22
# карточек подряд, 0 блоков). referer на detail-GET = URL avito-search страницы.
_AVITO_WARM_SEARCH_URL = "https://www.avito.ru/ekaterinburg/kvartiry/prodam?q=" + quote(
"купить квартиру"
)
_AVITO_WARM_YANDEX_URL = "https://yandex.ru/search/?text=" + quote(
"купить квартиру екатеринбург авито"
)
_AVITO_WARM_YANDEX_REFERER = "https://yandex.ru/"
# Reconnect-retry на detail-странице под backconnect-прокси — зеркало SERP-фикса # Reconnect-retry на detail-странице под backconnect-прокси — зеркало SERP-фикса
# (#1765/#1769, avito.py): 403/firewall = залочен лишь текущий exit-IP backconnect-прокси; # (#1765/#1769, avito.py): 403/firewall = залочен лишь текущий exit-IP backconnect-прокси;
# пересоздание curl_cffi-сессии = новый CONNECT-туннель = свежий exit-IP → escape. # пересоздание curl_cffi-сессии = новый CONNECT-туннель = свежий exit-IP → escape.
@ -283,6 +296,56 @@ def _build_detail_session() -> AsyncSession:
) )
async def warm_up_session(session: AsyncSession) -> bool:
"""Органический прогрев сессии: Yandex SERP (best-effort) -> Avito EKB search
(Referer=yandex). Сеет антибот-куки, после чего сессия держит батч detail-GET'ов
без 403. Возвращает True если search-страница загрузилась (не firewall).
"""
try:
await session.get(_AVITO_WARM_YANDEX_URL) # органический referer-источник, best-effort
except Exception as exc:
logger.debug("avito warm-up: yandex GET failed (non-fatal): %s", exc)
await asyncio.sleep(random.uniform(2.0, 3.5))
try:
r = await session.get(
_AVITO_WARM_SEARCH_URL, headers={"Referer": _AVITO_WARM_YANDEX_REFERER}
)
except Exception as exc:
logger.warning("avito warm-up: search fetch error: %s", exc)
return False
ok = r.status_code == 200 and not _is_firewall_page(r.text)
if not ok:
logger.warning("avito warm-up: search page blocked sc=%s", r.status_code)
return ok
async def research_in_session(session: AsyncSession) -> bool:
"""Лёгкий in-session перепоиск: re-GET avito search-страницы (Referer=yandex) на
ТЕКУЩЕЙ сессии освежает антибот-куки БЕЗ смены exit-IP. Best-effort (search-страница
иногда 429, это не ломает detail-фетчи). Возвращает True если страница не firewall.
"""
try:
r = await session.get(
_AVITO_WARM_SEARCH_URL, headers={"Referer": _AVITO_WARM_YANDEX_REFERER}
)
except Exception as exc:
logger.warning("avito re-search: fetch error: %s", exc)
return False
ok = r.status_code == 200 and not _is_firewall_page(r.text)
if not ok:
logger.info("avito re-search: search page sc=%s (non-fatal)", r.status_code)
return ok
async def build_warmed_session() -> AsyncSession:
"""Построить detail-сессию (_build_detail_session) и прогреть её. Возвращает
прогретую сессию (даже если warm-up soft-failed caller может всё равно пробовать).
"""
session = _build_detail_session()
await warm_up_session(session)
return session
# ── Soft-block detection (#1967 / #1950) ────────────────────────────────────── # ── Soft-block detection (#1967 / #1950) ──────────────────────────────────────
# Под нагрузкой Avito иногда вместо detail-страницы отдаёт HTTP 200 с ОБЩЕЙ # Под нагрузкой Avito иногда вместо detail-страницы отдаёт HTTP 200 с ОБЩЕЙ
# витриной (<title>Авито — Объявления на сайте Авито</title>) вместо запрошенного # витриной (<title>Авито — Объявления на сайте Авито</title>) вместо запрошенного
@ -327,11 +390,19 @@ async def fetch_detail(
*, *,
cffi_session: AsyncSession | None = None, cffi_session: AsyncSession | None = None,
browser_fetcher: BrowserFetcher | None = None, browser_fetcher: BrowserFetcher | None = None,
referer: str | None = None,
reconnect_on_block: bool = True,
) -> DetailEnrichment: ) -> DetailEnrichment:
"""GET <avito>/{item_url} → parse HTML via selectolax → DetailEnrichment. """GET <avito>/{item_url} → parse HTML via selectolax → DetailEnrichment.
Если browser_fetcher передан использует браузерный fetch (browser mode). Если browser_fetcher передан использует браузерный fetch (browser mode).
Если cffi_session не передана создаёт новую (impersonate='chrome120'). Если cffi_session не передана создаёт новую (impersonate='chrome120').
referer если задан, шлётся в Referer-заголовке detail-GET'а (warm-batch #1551:
referer = URL avito-search страницы, на которой прогрелась сессия).
reconnect_on_block если False, 403/firewall поднимает AvitoBlockedError СРАЗУ
(без холодных reconnect'ов — ими займётся outer-loop через rebuild+rewarm прогретой
сессии); 429 после short-retry поднимает AvitoRateLimitedError. Дефолт True
сохраняет старое поведение own-session / scrape_pipeline путей.
Raises: Raises:
httpx.HTTPError если status != 200 (через raise_for_status-like). httpx.HTTPError если status != 200 (через raise_for_status-like).
ValueError если item_id не извлечён из HTML. ValueError если item_id не извлечён из HTML.
@ -387,7 +458,12 @@ async def fetch_detail(
# ТОЛЬКО на наличие proxy (зеркало SERP-фикса #1765). Без proxy — старое поведение # ТОЛЬКО на наличие proxy (зеркало SERP-фикса #1765). Без proxy — старое поведение
# (raise на первом 403). Caller-provided сессию НЕ трогаем: reconnect использует # (raise на первом 403). Caller-provided сессию НЕ трогаем: reconnect использует
# эфемерную retry_session, которую сами закрываем в finally. # эфемерную retry_session, которую сами закрываем в finally.
backconnect = bool(settings.scraper_proxy_url) # reconnect_on_block=False (warm-batch #1551): не делаем холодных per-card
# reconnect'ов — блок поднимает AvitoBlockedError сразу, outer-loop пересоздаёт+
# прогревает сессию (свежий exit-IP). reconnect_on_block=True (own-session /
# scrape_pipeline) — старое поведение: backconnect-reconnect эфемерной сессией.
backconnect = bool(settings.scraper_proxy_url) and reconnect_on_block
req_headers = {"Referer": referer} if referer else None
r403 = 0 r403 = 0
r429 = 0 r429 = 0
r429_reconnect = 0 r429_reconnect = 0
@ -395,7 +471,7 @@ async def fetch_detail(
retry_session: AsyncSession | None = None retry_session: AsyncSession | None = None
try: try:
while True: while True:
response = await attempt_session.get(full_url) response = await attempt_session.get(full_url, headers=req_headers)
sc = response.status_code sc = response.status_code
is_firewall = sc == 200 and ( is_firewall = sc == 200 and (
_is_firewall_page(response.text) or _is_detail_soft_block(response.text) _is_firewall_page(response.text) or _is_detail_soft_block(response.text)

View file

@ -17,6 +17,7 @@ from __future__ import annotations
import asyncio import asyncio
import logging import logging
import random
import time import time
from dataclasses import dataclass, field from dataclasses import dataclass, field
from urllib.parse import urlparse from urllib.parse import urlparse
@ -30,7 +31,13 @@ from app.core.shutdown import shutdown_requested
from app.services import scrape_runs as runs_mod from app.services import scrape_runs as runs_mod
from app.services.scrape_pipeline import _CHROME_HEADERS, _avito_proxies from app.services.scrape_pipeline import _CHROME_HEADERS, _avito_proxies
from app.services.scrapers.avito import AvitoScraper from app.services.scrapers.avito import AvitoScraper
from app.services.scrapers.avito_detail import fetch_detail, save_detail_enrichment from app.services.scrapers.avito_detail import (
_AVITO_WARM_SEARCH_URL,
build_warmed_session,
fetch_detail,
research_in_session,
save_detail_enrichment,
)
from app.services.scrapers.avito_exceptions import ( from app.services.scrapers.avito_exceptions import (
AvitoBlockedError, AvitoBlockedError,
AvitoListingGoneError, AvitoListingGoneError,
@ -89,6 +96,10 @@ async def run_avito_detail_backfill(
budget_sec = float(params.get("budget_sec", 3600)) budget_sec = float(params.get("budget_sec", 3600))
request_delay_sec = float(params.get("request_delay_sec", 6.0)) request_delay_sec = float(params.get("request_delay_sec", 6.0))
max_consecutive_blocks = int(params.get("max_consecutive_blocks", 5)) max_consecutive_blocks = int(params.get("max_consecutive_blocks", 5))
warm_batch = int(params.get("warm_batch", 500))
research_every = int(params.get("research_every", 50))
block_cooldown_sec = float(params.get("block_cooldown_sec", 30.0))
warm_timeout_s = float(params.get("warm_timeout_s", 90.0))
# #1950: hard-timeout'ы на блокирующие await'ы внутри loop'а. Без них один # #1950: hard-timeout'ы на блокирующие await'ы внутри loop'а. Без них один
# зависший fetch_detail/rotate ронял весь run в zombie (run 423 завис 7.7ч). # зависший fetch_detail/rotate ронял весь run в zombie (run 423 завис 7.7ч).
@ -129,6 +140,14 @@ async def run_avito_detail_backfill(
await browser_fetcher.__aenter__() await browser_fetcher.__aenter__()
own_browser = True own_browser = True
scraper._browser = browser_fetcher scraper._browser = browser_fetcher
elif use_curl:
# use_curl=True (#1551 warm-batch): прогретая shared-сессия на sticky МГТС-IP —
# yandex-referer -> avito-search сеет антибот-куки, сессия держит батч detail без
# 403 (доказано пробами: ≥66 карточек подряд, 0 блоков). NB: МГТС sticky — один
# фикс. exit-IP, per-connection ротации нет (rebuild != новый IP); on-block —
# cooldown + in-session re-search, не дискард сессии.
own_session = True
session = await build_warmed_session()
elif not use_curl: elif not use_curl:
# curl_cffi legacy path (scraper_fetch_mode="curl_cffi", use_curl=False): # curl_cffi legacy path (scraper_fetch_mode="curl_cffi", use_curl=False):
# строим shared сессию через auv, как делал scrape_pipeline. # строим shared сессию через auv, как делал scrape_pipeline.
@ -140,8 +159,6 @@ async def run_avito_detail_backfill(
proxies=_avito_proxies(), proxies=_avito_proxies(),
) )
scraper._cffi = session scraper._cffi = session
# use_curl=True: ничего не строим — fetch_detail вызывает _build_detail_session()
# (scraper_proxy_url = backconnect mproxy) на каждый запрос.
runs_mod.update_heartbeat(db, run_id, current_counters) runs_mod.update_heartbeat(db, run_id, current_counters)
@ -197,6 +214,7 @@ async def run_avito_detail_backfill(
consecutive_blocks = 0 consecutive_blocks = 0
do_sleep = False do_sleep = False
items_since_warm = 0
for idx, row in enumerate(snapshot): for idx, row in enumerate(snapshot):
# Budget guard # Budget guard
@ -227,9 +245,10 @@ async def run_avito_detail_backfill(
) )
break break
# Delay before each request except the first # Delay before each request except the first (джиттер ±30% — менее
# роботизированный паттерн на органически прогретой сессии).
if do_sleep: if do_sleep:
await asyncio.sleep(request_delay_sec) await asyncio.sleep(request_delay_sec * random.uniform(0.7, 1.4))
do_sleep = True do_sleep = True
source_url: str = row["source_url"] source_url: str = row["source_url"]
@ -238,6 +257,52 @@ async def run_avito_detail_backfill(
# Normalise URL -> path (mirrors scrape_pipeline.py line 311) # Normalise URL -> path (mirrors scrape_pipeline.py line 311)
item_url = urlparse(source_url).path if source_url.startswith("http") else source_url item_url = urlparse(source_url).path if source_url.startswith("http") else source_url
# use_curl warm-batch (#1551): интервалы по размеру SERP-страницы (~50).
# МГТС sticky-IP: exit-IP константа (rebuild НЕ даёт новый IP). Каждые
# warm_batch (~500) — полный re-warm (close+build) на ТОМ ЖЕ sticky exit-IP:
# session-hygiene (свежий TLS-хендшейк + свежие куки yandex->search), НЕ смена
# IP. Между ними — лёгкий in-session перепоиск каждые research_every (~50)
# (re-GET search той же сессией, освежить куки). On-block — cooldown ниже.
if use_curl and session is not None:
if warm_batch and items_since_warm >= warm_batch:
# периодический полный re-warm на том же sticky exit-IP (свежий TLS+куки).
# #1950: bound сетевой re-warm; строим НОВУЮ сессию до close старой — на
# timeout/fail сохраняем текущую рабочую сессию (graceful degrade).
try:
new_session = await asyncio.wait_for(
build_warmed_session(), timeout=warm_timeout_s
)
try:
await session.close()
except Exception:
pass
session = new_session
except Exception:
logger.warning(
"avito_detail_backfill: run_id=%d re-warm timed out/failed (>%.0fs) "
"-- keeping current session",
run_id,
warm_timeout_s,
exc_info=True,
)
items_since_warm = 0
elif (
research_every
and items_since_warm > 0
and items_since_warm % research_every == 0
):
# лёгкий in-session перепоиск: тот же exit-IP, освежить куки.
# #1950: bound сетевой re-search.
try:
await asyncio.wait_for(research_in_session(session), timeout=warm_timeout_s)
except Exception:
logger.warning(
"avito_detail_backfill: run_id=%d in-session re-search "
"timed out/failed",
run_id,
exc_info=True,
)
try: try:
# #1950: hard-timeout — зависший fetch (browser hang / curl-stall) не # #1950: hard-timeout — зависший fetch (browser hang / curl-stall) не
# должен блокировать loop навсегда (иначе budget-guard/heartbeat молчат # должен блокировать loop навсегда (иначе budget-guard/heartbeat молчат
@ -247,11 +312,15 @@ async def run_avito_detail_backfill(
item_url, item_url,
cffi_session=session, cffi_session=session,
browser_fetcher=browser_fetcher, browser_fetcher=browser_fetcher,
referer=_AVITO_WARM_SEARCH_URL if use_curl else None,
reconnect_on_block=not use_curl,
), ),
timeout=fetch_timeout_s, timeout=fetch_timeout_s,
) )
if save_detail_enrichment(db, enrichment): if save_detail_enrichment(db, enrichment):
counters.enriched += 1 counters.enriched += 1
if use_curl:
items_since_warm += 1
consecutive_blocks = 0 consecutive_blocks = 0
except AvitoListingGoneError: except AvitoListingGoneError:
@ -299,18 +368,8 @@ async def run_avito_detail_backfill(
consecutive_blocks, consecutive_blocks,
e, e,
) )
# #1950: bound rotate — _rotate_ip ждёт settle + до 3 changeip-попыток; # abort FIRST — на последнем блоке не тратим cooldown+research/rotate
# на блокирующем changeip (зависшее соединение) без timeout loop виснет. # (consecutive_blocks уже инкрементнут выше).
try:
await asyncio.wait_for(scraper._rotate_ip(), timeout=rotate_timeout_s)
except Exception:
logger.warning(
"avito_detail_backfill: run_id=%d rotate_ip timed out/failed "
"(>%.0fs) -- continuing",
run_id,
rotate_timeout_s,
exc_info=True,
)
if consecutive_blocks >= max_consecutive_blocks: if consecutive_blocks >= max_consecutive_blocks:
logger.error( logger.error(
"avito_detail_backfill: run_id=%d ABORT -- %d consecutive blocks, " "avito_detail_backfill: run_id=%d ABORT -- %d consecutive blocks, "
@ -321,6 +380,40 @@ async def run_avito_detail_backfill(
counters.attempted, counters.attempted,
) )
break break
# МГТС sticky-IP: один фикс. exit-IP, per-connection ротации нет (проверено:
# 6/6 свежих сессий = тот же IP 109.252.125.80; ротация только вручную
# кнопкой). Уйти на свежий IP софтом нельзя → блок = rate-limit текущего IP:
# даём окну остыть (cooldown) и освежаем куки in-session, НЕ дискардим
# рабочую прогретую сессию (rebuild на том же IP бесполезен для escape +
# грузит rate-limited IP полным прогревом). items_since_warm НЕ сбрасываем.
if use_curl:
await asyncio.sleep(block_cooldown_sec)
if session is not None:
# #1950: bound сетевой re-search.
try:
await asyncio.wait_for(
research_in_session(session), timeout=warm_timeout_s
)
except Exception:
logger.warning(
"avito_detail_backfill: run_id=%d on-block re-search "
"timed out/failed",
run_id,
exc_info=True,
)
else:
# #1950: bound rotate — _rotate_ip ждёт settle + до 3 changeip-попыток;
# на блокирующем changeip (зависшее соединение) без timeout loop виснет.
try:
await asyncio.wait_for(scraper._rotate_ip(), timeout=rotate_timeout_s)
except Exception:
logger.warning(
"avito_detail_backfill: run_id=%d rotate_ip timed out/failed "
"(>%.0fs) -- continuing",
run_id,
rotate_timeout_s,
exc_info=True,
)
except TimeoutError: except TimeoutError:
# asyncio.wait_for → TimeoutError (py3.12: asyncio.TimeoutError — alias). # asyncio.wait_for → TimeoutError (py3.12: asyncio.TimeoutError — alias).

View file

@ -0,0 +1,140 @@
"""Живой прогрев сессии + батч-на-сессию (#1551).
warm_up_session: GET Yandex SERP (best-effort) -> GET Avito EKB search-страница с
Referer=yandex сеет антибот-куки. Прогретая сессия потом держит батч detail-GET'ов
без 403 (доказано пробами на проде: 22/22 карточек подряд, 0 блоков).
fetch_detail(referer=..., reconnect_on_block=False): шлёт Referer detail-GET'а и при
403/firewall поднимает AvitoBlockedError СРАЗУ без холодных per-card reconnect'ов
(их делает outer-loop через rebuild+rewarm прогретой сессии в avito_detail_backfill).
Стиль фикстур/моков зеркалит test_avito_detail_403_reconnect.py.
"""
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
from app.core.config import settings
from app.services.scrapers import avito_detail as detail_mod
from app.services.scrapers.avito_detail import (
fetch_detail,
research_in_session,
warm_up_session,
)
from app.services.scrapers.avito_exceptions import AvitoBlockedError
def _resp(status_code: int, text: str = "") -> MagicMock:
r = MagicMock()
r.status_code = status_code
r.text = text
return r
def _enable_backconnect(monkeypatch: pytest.MonkeyPatch) -> None:
monkeypatch.setattr(settings, "scraper_proxy_url_env", "http://u:p@mproxy.site:14619")
# ── warm_up_session ───────────────────────────────────────────────────────────
@pytest.mark.asyncio
async def test_warm_up_session_search_with_yandex_referer_true_on_200() -> None:
"""warm_up_session шлёт GET на _AVITO_WARM_SEARCH_URL с Referer=yandex и
возвращает True на 200/non-firewall."""
session = AsyncMock()
session.get = AsyncMock(return_value=_resp(200, "<html>ok</html>"))
with (
patch("app.services.scrapers.avito_detail.asyncio.sleep", AsyncMock()),
patch("app.services.scrapers.avito_detail._is_firewall_page", return_value=False),
):
ok = await warm_up_session(session)
assert ok is True
# search-GET (на avito-search URL) ушёл ровно один раз и с Referer=yandex.
search_calls = [
c
for c in session.get.await_args_list
if c.args and c.args[0] == detail_mod._AVITO_WARM_SEARCH_URL
]
assert len(search_calls) == 1
assert search_calls[0].kwargs["headers"]["Referer"] == detail_mod._AVITO_WARM_YANDEX_REFERER
# Yandex SERP — органический referer-источник — тоже фетчился (best-effort, до search).
assert any(
c.args and c.args[0] == detail_mod._AVITO_WARM_YANDEX_URL
for c in session.get.await_args_list
)
@pytest.mark.asyncio
async def test_warm_up_session_false_on_403() -> None:
"""search-страница вернула 403 → warm_up_session возвращает False."""
session = AsyncMock()
session.get = AsyncMock(return_value=_resp(403))
with (
patch("app.services.scrapers.avito_detail.asyncio.sleep", AsyncMock()),
patch("app.services.scrapers.avito_detail._is_firewall_page", return_value=False),
):
ok = await warm_up_session(session)
assert ok is False
# ── fetch_detail referer + reconnect_on_block=False ───────────────────────────
@pytest.mark.asyncio
async def test_fetch_detail_referer_forwarded_and_block_raises_immediately(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""fetch_detail(referer="X", reconnect_on_block=False): Referer="X" уходит в
detail-GET, а 403 поднимает AvitoBlockedError СРАЗУ без холодного reconnect
(даже когда backconnect-proxy включён гейт = reconnect_on_block, не наличие proxy).
"""
_enable_backconnect(monkeypatch)
caller = AsyncMock()
caller.get = AsyncMock(return_value=_resp(403))
build = MagicMock() # _build_detail_session НЕ должна дёргаться (нет reconnect'а)
with patch.object(detail_mod, "_build_detail_session", build):
with pytest.raises(AvitoBlockedError, match="HTTP 403"):
await fetch_detail(
"/ekaterinburg/kvartiry/test-123",
cffi_session=caller,
referer="X",
reconnect_on_block=False,
)
build.assert_not_called() # никаких холодных reconnect'ов
assert caller.get.await_count == 1 # один detail-GET → сразу raise
assert caller.get.await_args.kwargs["headers"] == {"Referer": "X"}
# ── research_in_session (in-session перепоиск, #1551 follow-up) ────────────────
@pytest.mark.asyncio
async def test_research_in_session_re_get_search_with_yandex_referer_true_on_200() -> None:
"""research_in_session re-GET'ит avito-search на ТОЙ ЖЕ сессии с Referer=yandex
и возвращает True на 200/non-firewall (свежесть куки без смены exit-IP)."""
session = AsyncMock()
session.get = AsyncMock(return_value=_resp(200, "<html>ok</html>"))
with patch("app.services.scrapers.avito_detail._is_firewall_page", return_value=False):
ok = await research_in_session(session)
assert ok is True
assert session.get.await_count == 1 # ровно один re-GET, без yandex-фазы (in-session)
assert session.get.await_args.args[0] == detail_mod._AVITO_WARM_SEARCH_URL
sent_referer = session.get.await_args.kwargs["headers"]["Referer"]
assert sent_referer == detail_mod._AVITO_WARM_YANDEX_REFERER
@pytest.mark.asyncio
async def test_research_in_session_false_on_403() -> None:
"""search вернул 403 → research_in_session возвращает False (best-effort, не fatal)."""
session = AsyncMock()
session.get = AsyncMock(return_value=_resp(403))
with patch("app.services.scrapers.avito_detail._is_firewall_page", return_value=False):
ok = await research_in_session(session)
assert ok is False

View file

@ -27,6 +27,25 @@ def _reset_shutdown() -> None:
_sd.reset_shutdown() _sd.reset_shutdown()
@pytest.fixture(autouse=True)
def _patch_build_warmed() -> object:
"""use_curl-ветка (#1551) строит прогретую сессию через build_warmed_session и
освежает куки через research_in_session (реальный yandex/avito GET + sleep).
fetch_detail в этих тестах всё равно замокан мокаем обе не-сетевыми AsyncMock,
чтобы setup / periodic re-warm / on-block cooldown не ходили в сеть."""
with (
patch(
"app.tasks.avito_detail_backfill.build_warmed_session",
AsyncMock(return_value=AsyncMock()),
),
patch(
"app.tasks.avito_detail_backfill.research_in_session",
AsyncMock(return_value=True),
) as research,
):
yield research
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Helpers # Helpers
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
@ -53,6 +72,8 @@ _SESSION = "app.tasks.avito_detail_backfill.AsyncSession"
_SCRAPER = "app.tasks.avito_detail_backfill.AvitoScraper" _SCRAPER = "app.tasks.avito_detail_backfill.AvitoScraper"
_SETTINGS = "app.tasks.avito_detail_backfill.settings" _SETTINGS = "app.tasks.avito_detail_backfill.settings"
_SHUTDOWN = "app.tasks.avito_detail_backfill.shutdown_requested" _SHUTDOWN = "app.tasks.avito_detail_backfill.shutdown_requested"
_BUILD_WARM = "app.tasks.avito_detail_backfill.build_warmed_session"
_RESEARCH = "app.tasks.avito_detail_backfill.research_in_session"
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Tests # Tests
@ -119,7 +140,11 @@ async def test_backfill_processes_snapshot_to_completion() -> None:
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_backfill_blocked_abort_after_max_consecutive() -> None: async def test_backfill_blocked_abort_after_max_consecutive() -> None:
"""5 consecutive AvitoBlockedError -> abort, mark_done (NOT mark_failed), rotate_ip x5.""" """5 consecutive AvitoBlockedError -> abort, mark_done (NOT mark_failed).
#1950 abort-reorder: abort-check ПЕРЕД recovery → на 5-м (аборт-)блоке rotate_ip
НЕ дёргается (не тратим recovery на финальном блоке). rotate_ip x4 (блоки 1-4).
"""
from app.services.scrapers.avito_exceptions import AvitoBlockedError from app.services.scrapers.avito_exceptions import AvitoBlockedError
snapshot = _make_snapshot(10) snapshot = _make_snapshot(10)
@ -129,7 +154,9 @@ async def test_backfill_blocked_abort_after_max_consecutive() -> None:
mock_fetch = AsyncMock(side_effect=blocked_exc) mock_fetch = AsyncMock(side_effect=blocked_exc)
mock_scraper = MagicMock() mock_scraper = MagicMock()
mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True) mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True)
fake_settings = MagicMock(scraper_fetch_mode="cffi") # use_curl=False: legacy block→_rotate_ip путь (#1551 warm-batch rebuild только при
# use_curl=True). MagicMock без явного флага сделал бы use_curl truthy.
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False)
with ( with (
patch(_SETTINGS, fake_settings), patch(_SETTINGS, fake_settings),
patch(_SESSION, return_value=AsyncMock()), patch(_SESSION, return_value=AsyncMock()),
@ -147,7 +174,8 @@ async def test_backfill_blocked_abort_after_max_consecutive() -> None:
assert result.enriched == 0 assert result.enriched == 0
runs.mark_done.assert_called_once() runs.mark_done.assert_called_once()
runs.mark_failed.assert_not_called() runs.mark_failed.assert_not_called()
assert mock_scraper.return_value._rotate_ip.call_count == 5 # abort-check до recovery → 5-й блок абортит без rotate; rotate только на блоках 1-4.
assert mock_scraper.return_value._rotate_ip.call_count == 4
@pytest.mark.asyncio @pytest.mark.asyncio
@ -249,7 +277,9 @@ async def test_backfill_rotate_ip_called_on_each_block() -> None:
mock_fetch = AsyncMock(side_effect=[blocked_exc, mock_enrichment]) mock_fetch = AsyncMock(side_effect=[blocked_exc, mock_enrichment])
mock_scraper = MagicMock() mock_scraper = MagicMock()
mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True) mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True)
fake_settings = MagicMock(scraper_fetch_mode="cffi") # use_curl=False: legacy block→_rotate_ip путь (#1551 warm-batch rebuild только при
# use_curl=True). MagicMock без явного флага сделал бы use_curl truthy.
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False)
with ( with (
patch(_SETTINGS, fake_settings), patch(_SETTINGS, fake_settings),
patch(_SESSION, return_value=AsyncMock()), patch(_SESSION, return_value=AsyncMock()),
@ -342,7 +372,14 @@ async def test_backfill_fetch_timeout_skips_and_continues() -> None:
call_urls: list[str] = [] call_urls: list[str] = []
async def _fetch(url: str, *, cffi_session: object = None, browser_fetcher: object = None): async def _fetch(
url: str,
*,
cffi_session: object = None,
browser_fetcher: object = None,
referer: object = None,
reconnect_on_block: bool = True,
):
call_urls.append(url) call_urls.append(url)
if len(call_urls) == 1: if len(call_urls) == 1:
await asyncio.Event().wait() # висит вечно → wait_for отменит по timeout await asyncio.Event().wait() # висит вечно → wait_for отменит по timeout
@ -394,7 +431,9 @@ async def test_backfill_listing_gone_marks_inactive_no_breaker() -> None:
mock_fetch = AsyncMock(side_effect=AvitoListingGoneError("404 gone")) mock_fetch = AsyncMock(side_effect=AvitoListingGoneError("404 gone"))
mock_scraper = MagicMock() mock_scraper = MagicMock()
mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True) mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True)
fake_settings = MagicMock(scraper_fetch_mode="cffi") # use_curl=False: legacy block→_rotate_ip путь (#1551 warm-batch rebuild только при
# use_curl=True). MagicMock без явного флага сделал бы use_curl truthy.
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False)
with ( with (
patch(_SETTINGS, fake_settings), patch(_SETTINGS, fake_settings),
patch(_SESSION, return_value=AsyncMock()), patch(_SESSION, return_value=AsyncMock()),
@ -510,3 +549,60 @@ async def test_backfill_use_curl_false_creates_browser_fetcher() -> None:
assert kwargs.get("browser_fetcher") is not None assert kwargs.get("browser_fetcher") is not None
assert result.enriched == 1 assert result.enriched == 1
runs.mark_done.assert_called_once() runs.mark_done.assert_called_once()
@pytest.mark.asyncio
async def test_backfill_use_curl_block_cooldown_research_no_rebuild() -> None:
"""#1551 sticky-IP: on-block (use_curl) даёт cooldown + in-session research, НЕ
пересоздаёт прогретую сессию и НЕ дёргает changeip-ротацию.
МГТС sticky один фикс. exit-IP (rebuild != новый IP), уйти на свежий IP софтом
нельзя. Блок = rate-limit текущего IP cooldown + re-search той же сессией.
fetch: 1-й вызов BLOCKED, 2-й успех.
"""
from app.services.scrapers.avito_exceptions import AvitoBlockedError
snapshot = _make_snapshot(2)
db = _mock_db(snapshot)
runs = MagicMock()
mock_enrichment = MagicMock()
blocked_exc = AvitoBlockedError("rate-limited")
mock_fetch = AsyncMock(side_effect=[blocked_exc, mock_enrichment])
warmed_session = AsyncMock()
fake_settings = MagicMock(
scraper_fetch_mode="cffi",
avito_detail_backfill_use_curl=True,
)
with (
patch(_SETTINGS, fake_settings),
patch(_SCRAPER) as mock_scraper,
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
patch(_BUILD_WARM, AsyncMock(return_value=warmed_session)) as mock_build,
patch(_RESEARCH, new_callable=AsyncMock) as mock_research,
):
result = await run_avito_detail_backfill(
db,
run_id=14,
params={
"batch_size": 10,
"budget_sec": 3600,
"max_consecutive_blocks": 5,
"block_cooldown_sec": 0.0,
},
)
assert result.blocked == 1
assert result.enriched == 1
assert result.attempted == 2
# on-block: research_in_session вызван (освежить куки in-session)...
mock_research.assert_awaited_once()
# ...а build_warmed_session НЕ дёргался повторно (вызван только 1× в setup —
# прогретая сессия НЕ пересоздаётся на блоке, sticky-IP rebuild бесполезен).
assert mock_build.await_count == 1
# changeip-ротация (legacy путь) под use_curl НЕ дёргается.
mock_scraper.return_value._rotate_ip.assert_not_called()
runs.mark_done.assert_called_once()
runs.mark_failed.assert_not_called()