from __future__ import annotations import asyncio import fnmatch import os import sys from typing import Any from unittest.mock import AsyncMock, MagicMock, patch os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") _wp_mock = MagicMock() sys.modules.setdefault("weasyprint", _wp_mock) import pytest # noqa: E402 from scraper_kit.avito_exceptions import ( # noqa: E402 AvitoBlockedError, AvitoSidecarUnavailableError, ) from app.core import shutdown as _sd # noqa: E402 from app.services.proxy_rotation import RotationResult # noqa: E402 from app.tasks.avito_detail_backfill import ( # noqa: E402 _DEFAULT_ROTATE_RECONNECT_DELAY_S, _OBLAST_AVITO_URL_PATTERNS, AvitoDetailBackfillResult, run_avito_detail_backfill, ) @pytest.fixture(autouse=True) def _reset_shutdown() -> None: """shutdown — module-global Event: чистим вокруг каждого теста (изоляция #1182).""" _sd.reset_shutdown() yield _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 # --------------------------------------------------------------------------- def _make_snapshot(n: int) -> list[dict]: return [{"id": i + 1, "source_url": f"/items/{i + 1}"} for i in range(n)] def _fake_settings(**overrides: object) -> MagicMock: """settings.* мокается атрибутами; bare MagicMock() без явных detail_backfill_block_ratio_window/_threshold отдаёт int(MagicMock())==1 -- окно=1, порог=1.0 абортит на первом же блоке (#3184 review MINOR 4). Единая точка дефолтов вместо 23 копий тех же двух kwargs на каждом call site.""" defaults: dict[str, object] = { "detail_backfill_block_ratio_window": 20, "detail_backfill_block_ratio_threshold": 0.7, # Реальный дефолт (config.py) -- без него bare MagicMock() возвращает # child-MagicMock на сравнение `attempts_since_rotation >= settings.avito_ # detail_backfill_rotate_after_attempts` в browser-режиме и падает TypeError # (int >= MagicMock не поддерживается). В curl-режиме (browser_fetcher=None) # сравнение вообще не вычисляется (short-circuit `and`), но browser-тесты # ниже (test_backfill_use_curl_false_creates_browser_fetcher и rotate-тесты) # его достигают. "avito_detail_backfill_rotate_after_attempts": 15, } defaults.update(overrides) return MagicMock(**defaults) def _mock_db(snapshot: list[dict]) -> MagicMock: """Fake Session: first execute() returns snapshot via .mappings().all().""" db = MagicMock() sel = MagicMock() sel.mappings.return_value.all.return_value = snapshot db.execute.return_value = sel return db _FETCH = "app.tasks.avito_detail_backfill.fetch_detail" _SAVE = "app.tasks.avito_detail_backfill.save_detail_enrichment" _RUNS = "app.tasks.avito_detail_backfill.runs_mod" _SLEEP = "app.tasks.avito_detail_backfill.asyncio.sleep" _SESSION = "app.tasks.avito_detail_backfill.AsyncSession" _SCRAPER = "app.tasks.avito_detail_backfill.AvitoScraper" _SETTINGS = "app.tasks.avito_detail_backfill.settings" _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" # #2825: settings.scraper_proxy_url в "elif not use_curl" (legacy curl_cffi) branch # заменён на resolve_proxy_url(db, "avito") (пул scrape_proxies с учётом банов, # fallback на settings.scraper_proxy_url внутри app.services.proxy_egress) -- эти # тесты про block/ban/rotate-логику, не про подбор прокси (см. # tests/services/test_proxy_egress.py), поэтому мокаем сам резолвер. _RESOLVE_PROXY_URL = "app.tasks.avito_detail_backfill.resolve_proxy_url" _ROTATE_PROXY = "app.tasks.avito_detail_backfill.rotate_proxy" # --------------------------------------------------------------------------- # Tests # --------------------------------------------------------------------------- @pytest.mark.asyncio async def test_backfill_empty_snapshot_marks_done() -> None: """Empty snapshot -> mark_done immediately, no fetch calls.""" db = _mock_db([]) runs = MagicMock() fake_session = AsyncMock() fake_settings = _fake_settings( scraper_fetch_mode="cffi", ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=fake_session), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH) as mock_fetch, ): result = await run_avito_detail_backfill( db, run_id=1, params={"batch_size": 10, "budget_sec": 60} ) assert isinstance(result, AvitoDetailBackfillResult) assert result.attempted == 0 assert result.enriched == 0 mock_fetch.assert_not_called() runs.mark_done.assert_called_once() runs.mark_failed.assert_not_called() @pytest.mark.asyncio async def test_backfill_processes_snapshot_to_completion() -> None: """3 listings -> all fetched and enriched, mark_done called.""" from app.services.scraper_adapters import RealScraperConfig snapshot = _make_snapshot(3) db = _mock_db(snapshot) runs = MagicMock() mock_enrichment = MagicMock() mock_fetch = AsyncMock(return_value=mock_enrichment) mock_save = MagicMock(return_value=True) fake_settings = _fake_settings( scraper_fetch_mode="cffi", ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER) as mock_scraper_cls, patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, mock_save), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=2, params={"batch_size": 10, "budget_sec": 3600} ) assert result.attempted == 3 assert result.enriched == 3 assert result.blocked == 0 assert result.failed == 0 assert mock_fetch.call_count == 3 runs.mark_backfill_finished.assert_called_once() runs.mark_failed.assert_not_called() # #2310 regression guard: kit fetch_detail silently drops the backconnect- # on-403 retry (and kit AvitoScraper can't read scraper_proxy_url at all) # unless config=RealScraperConfig() is passed/injected at the call # site — assert_called()/call_count alone wouldn't catch someone dropping # that kwarg later (mirrors #2306's test_backfill_wave2.py:282-286 pattern). _, fetch_call_kwargs = mock_fetch.call_args assert isinstance(fetch_call_kwargs.get("config"), RealScraperConfig) assert isinstance(mock_scraper_cls.call_args.args[0], RealScraperConfig) @pytest.mark.asyncio async def test_backfill_build_warmed_session_receives_config() -> None: """#2397 Part D1 / #2330 regression guard: kit's build_warmed_session() now accepts config=ScraperConfig and silently drops settings.scraper_proxy_url (sticky МГТС-прокси) -- NO crash, just a proxy-less session -- unless config=RealScraperConfig() is passed at the call site. use_curl=True is the prod-default path (avito_detail_backfill_use_curl), so this is the setup-time build_warmed_session() call (~line 177). assert_awaited_once() alone wouldn't catch someone dropping the kwarg later (mirrors the fetch_detail config= guard above, lines 141-148).""" from app.services.scraper_adapters import RealScraperConfig snapshot = _make_snapshot(1) db = _mock_db(snapshot) runs = MagicMock() mock_enrichment = MagicMock() mock_fetch = AsyncMock(return_value=mock_enrichment) mock_save = MagicMock(return_value=True) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, mock_save), patch(_SLEEP, new_callable=AsyncMock), patch(_BUILD_WARM, AsyncMock(return_value=AsyncMock())) as mock_build, ): result = await run_avito_detail_backfill( db, run_id=15, params={"batch_size": 10, "budget_sec": 3600} ) assert result.enriched == 1 mock_build.assert_awaited_once() _, build_call_kwargs = mock_build.call_args assert isinstance(build_call_kwargs.get("config"), RealScraperConfig) runs.mark_backfill_finished.assert_called_once() runs.mark_failed.assert_not_called() @pytest.mark.asyncio @pytest.mark.parametrize( ("exc_factory", "expected_kinds"), [ (lambda: AvitoBlockedError("ip blocked"), {"platform"}), (lambda: AvitoSidecarUnavailableError("browser unavailable"), {"infra"}), ], ) async def test_backfill_reports_ban_kind_of_the_blocks_it_saw( exc_factory: Any, expected_kinds: set[str] ) -> None: """Диагноз блоков доезжает до финализатора по ТИПУ исключения (#2764). Фальсификация: до правки задача не передавала ничего, и обе серии — отказ площадки и отказ нашего сайдкара — давали в scrape_runs.ban_kind одинаковое 'platform' по умолчанию (прод, прогон 3306: blocked=5, причина не установлена). """ snapshot = _make_snapshot(10) db = _mock_db(snapshot) runs = MagicMock() mock_scraper = MagicMock() mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER, mock_scraper), patch(_RUNS, runs), patch(_RESOLVE_PROXY_URL, MagicMock(return_value="http://test-proxy.local:8080")), patch(_FETCH, AsyncMock(side_effect=exc_factory())), patch(_SLEEP, new_callable=AsyncMock), ): await run_avito_detail_backfill( db, run_id=3, params={"batch_size": 10, "budget_sec": 3600, "max_consecutive_blocks": 5} ) runs.mark_backfill_finished.assert_called_once() assert set(runs.mark_backfill_finished.call_args.kwargs["ban_kinds"]) == expected_kinds @pytest.mark.asyncio async def test_backfill_ban_kinds_census_keeps_multiplicities() -> None: """4×AvitoBlockedError + 1×AvitoSidecarUnavailableError -> перепись 4/1, не set (#3178). Фальсификация: до правки block_ban_kinds был set() -- .add() схлопнул бы этот прогон в тот же {'platform', 'infra'}, что и настоящие 2+2, и mark_backfill_finished не смог бы отличить явное большинство от честной ничьей. """ from scraper_kit.avito_exceptions import AvitoBlockedError, AvitoSidecarUnavailableError snapshot = _make_snapshot(10) db = _mock_db(snapshot) runs = MagicMock() # abort-check срабатывает ПЕРЕД recovery на 5-м блоке -- ровно 5 попыток. mock_fetch = AsyncMock( side_effect=[ AvitoBlockedError("firewall"), AvitoBlockedError("firewall"), AvitoBlockedError("firewall"), AvitoBlockedError("firewall"), AvitoSidecarUnavailableError("browser unavailable"), ] ) mock_scraper = MagicMock() mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER, mock_scraper), patch(_RUNS, runs), patch(_RESOLVE_PROXY_URL, MagicMock(return_value="http://test-proxy.local:8080")), patch(_FETCH, mock_fetch), patch(_SLEEP, new_callable=AsyncMock), ): await run_avito_detail_backfill( db, run_id=3, params={"batch_size": 10, "budget_sec": 3600, "max_consecutive_blocks": 5} ) runs.mark_backfill_finished.assert_called_once() census = runs.mark_backfill_finished.call_args.kwargs["ban_kinds"] assert dict(census) == {"platform": 4, "infra": 1} @pytest.mark.asyncio async def test_backfill_abort_log_has_no_ip_rate_limited_literal(caplog: Any) -> None: """ABORT-лог называет измеренную причину, а не литерал 'IP rate-limited' (#3178). До правки текст ABORT всегда писал 'IP rate-limited' -- даже когда блоки были отказом НАШЕГО сайдкара (AvitoSidecarUnavailableError), не площадки. """ from scraper_kit.avito_exceptions import AvitoBlockedError snapshot = _make_snapshot(10) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(side_effect=AvitoBlockedError("firewall/soft-block")) mock_scraper = MagicMock() mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER, mock_scraper), patch(_RUNS, runs), patch(_RESOLVE_PROXY_URL, MagicMock(return_value="http://test-proxy.local:8080")), patch(_FETCH, mock_fetch), patch(_SLEEP, new_callable=AsyncMock), caplog.at_level("ERROR"), ): await run_avito_detail_backfill( db, run_id=3, params={"batch_size": 10, "budget_sec": 3600, "max_consecutive_blocks": 5} ) abort_records = [r.message for r in caplog.records if "ABORT" in r.message] assert abort_records, "ожидался ABORT-лог" assert "IP rate-limited" not in abort_records[0] assert "AvitoBlockedError" in abort_records[0] @pytest.mark.asyncio async def test_backfill_blocked_abort_after_max_consecutive() -> None: """5 consecutive AvitoBlockedError -> abort с пометкой aborted_by_blocks (#2674). Раньше — mark_done; на проде 13 прогонов attempted=5 blocked=5 enriched=0 назывались успехом. Теперь флаг обрыва → статус 'banned'. #1950 abort-reorder: abort-check ПЕРЕД recovery → на 5-м (аборт-)блоке rotate_ip НЕ дёргается (не тратим recovery на финальном блоке). rotate_ip x4 (блоки 1-4). """ from scraper_kit.avito_exceptions import AvitoBlockedError snapshot = _make_snapshot(10) db = _mock_db(snapshot) runs = MagicMock() blocked_exc = AvitoBlockedError("ip blocked") mock_fetch = AsyncMock(side_effect=blocked_exc) mock_scraper = MagicMock() mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True) # use_curl=False: legacy block→_rotate_ip путь (#1551 warm-batch rebuild только при # use_curl=True). MagicMock без явного флага сделал бы use_curl truthy. fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER, mock_scraper), patch(_RUNS, runs), patch(_RESOLVE_PROXY_URL, MagicMock(return_value="http://test-proxy.local:8080")), patch(_FETCH, mock_fetch), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=3, params={"batch_size": 10, "budget_sec": 3600, "max_consecutive_blocks": 5} ) assert result.blocked == 5 assert result.attempted == 5 assert result.enriched == 0 runs.mark_backfill_finished.assert_called_once() assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is True runs.mark_failed.assert_not_called() # abort-check до recovery → 5-й блок абортит без rotate; rotate только на блоках 1-4. assert mock_scraper.return_value._rotate_ip.call_count == 4 @pytest.mark.asyncio async def test_backfill_sigterm_drain_breaks_and_marks_done_partial() -> None: """#1182 Phase 2: shutdown_requested() True → loop выходит на границе карточки, mark_done вызывается с ЧАСТИЧНЫМИ счётчиками (не mark_failed, не mark_cancelled). side_effect [False, True]: 1-я карточка обрабатывается (attempted=1, enriched=1), перед 2-й приходит SIGTERM-drain → break. Snapshot — pending-query, остаток до-резюмит следующий run (resume_cursor не нужен). """ snapshot = _make_snapshot(3) db = _mock_db(snapshot) runs = MagicMock() mock_enrichment = MagicMock() mock_fetch = AsyncMock(return_value=mock_enrichment) mock_save = MagicMock(return_value=True) fake_settings = _fake_settings( scraper_fetch_mode="cffi", ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, mock_save), patch(_SLEEP, new_callable=AsyncMock), patch(_SHUTDOWN, side_effect=[False, True]), ): result = await run_avito_detail_backfill( db, run_id=13, params={"batch_size": 10, "budget_sec": 3600} ) # Обработана только 1-я карточка, на 2-й — drain-break. assert result.attempted == 1 assert result.enriched == 1 assert mock_fetch.call_count == 1 runs.mark_backfill_finished.assert_called_once() runs.mark_failed.assert_not_called() runs.mark_cancelled.assert_not_called() # mark_done получил ЧАСТИЧНЫЕ счётчики (attempted=1, а не весь snapshot=3). done_counters = runs.mark_backfill_finished.call_args.args[2] assert done_counters["attempted"] == 1 @pytest.mark.asyncio async def test_backfill_budget_guard_stops_loop() -> None: """Budget expired before first listing -> fetch_detail not called.""" snapshot = _make_snapshot(5) db = _mock_db(snapshot) runs = MagicMock() mono_values = iter([0.0, 999.0, 999.0]) fake_settings = _fake_settings( scraper_fetch_mode="cffi", ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch("app.tasks.avito_detail_backfill.time.monotonic", side_effect=mono_values), patch(_FETCH) as mock_fetch, ): await run_avito_detail_backfill(db, run_id=4, params={"batch_size": 5, "budget_sec": 1}) mock_fetch.assert_not_called() runs.mark_backfill_finished.assert_called_once() @pytest.mark.asyncio async def test_backfill_top_level_exception_marks_failed() -> None: """db.execute raises -> mark_failed called, exception re-raised.""" db = MagicMock() db.execute.side_effect = RuntimeError("DB connection lost") runs = MagicMock() fake_settings = _fake_settings( scraper_fetch_mode="cffi", ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), ): with pytest.raises(RuntimeError, match="DB connection lost"): await run_avito_detail_backfill( db, run_id=5, params={"batch_size": 5, "budget_sec": 60} ) runs.mark_failed.assert_called_once() runs.mark_backfill_finished.assert_not_called() @pytest.mark.asyncio async def test_backfill_rotate_ip_called_on_each_block() -> None: """1 block + 1 success -> rotate_ip called once, enriched=1.""" from scraper_kit.avito_exceptions import AvitoBlockedError snapshot = _make_snapshot(2) db = _mock_db(snapshot) runs = MagicMock() mock_enrichment = MagicMock() blocked_exc = AvitoBlockedError("x") mock_fetch = AsyncMock(side_effect=[blocked_exc, mock_enrichment]) mock_scraper = MagicMock() mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True) # use_curl=False: legacy block→_rotate_ip путь (#1551 warm-batch rebuild только при # use_curl=True). MagicMock без явного флага сделал бы use_curl truthy. fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER, mock_scraper), patch(_RUNS, runs), patch(_RESOLVE_PROXY_URL, MagicMock(return_value="http://test-proxy.local:8080")), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=6, params={"batch_size": 10, "budget_sec": 3600, "max_consecutive_blocks": 5} ) assert result.enriched == 1 assert result.blocked == 1 assert mock_scraper.return_value._rotate_ip.call_count == 1 runs.mark_backfill_finished.assert_called_once() @pytest.mark.asyncio async def test_backfill_snapshot_filters_ekb_active_only() -> None: """Снапшот-SELECT (#1814, расширено #2576) фильтрует активные ЕКБ- И известные oblast-листинги (region 66), НЕ всё подряд. Проверяем, что текст запроса содержит `is_active = TRUE`, `LIKE '%/ekaterinburg/%'` (ekb CTE, LIMIT batch_size НЕ сокращён) и `LIKE ANY(...)` по oblast-паттернам (oblast CTE, отдельный LIMIT oblast_batch_size) — legacy не-ЕКБ/не-область (moskva/spb/tyumen) и мёртвые листинги не попадают в фетч, иначе browser спотыкается → curl-бан 429. """ db = _mock_db([]) runs = MagicMock() fake_settings = _fake_settings( scraper_fetch_mode="cffi", ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH), ): await run_avito_detail_backfill(db, run_id=8, params={"batch_size": 10, "budget_sec": 60}) # Первый (и единственный при пустом снапшоте) execute — это SELECT-снапшот. snapshot_call = db.execute.call_args_list[0] sql_text = str(snapshot_call.args[0]) assert "is_active = TRUE" in sql_text assert "/ekaterinburg/" in sql_text assert "LIKE ANY(CAST(:oblast_patterns AS text[]))" in sql_text assert "detail_enriched_at IS NULL" in sql_text assert "(lat IS NULL) DESC" in sql_text assert "CAST(:batch_size AS int)" in sql_text assert "CAST(:oblast_batch_size AS int)" in sql_text # ekb-квота передаётся batch_size БЕЗ урезания (#2576 требование "ЕКБ не # деградирует") — oblast получает отдельный (не вычтенный) bind-параметр. bind_params = snapshot_call.args[1] assert bind_params["batch_size"] == 10 assert bind_params["oblast_batch_size"] == 100 # default assert set(bind_params["oblast_patterns"]) == set(_OBLAST_AVITO_URL_PATTERNS) def test_oblast_avito_url_patterns_cover_region66_cities() -> None: """#2576: _OBLAST_AVITO_URL_PATTERNS строится из CITY_LOCATIONS.avito_slug — список должен покрывать реальные Avito-слаги oblast-городов (в т.ч. те, что ОТЛИЧАЮТСЯ от нашего city_slug: kamensk-uralskiy через дефис, а не kamensk_uralskiy). #2578 review: '_' в слаге -- LIKE wildcard, экранируем при построении паттерна ('_' -> '\\_') -- nizhniy_tagil/verhnyaya_pyshma здесь ожидаются С обратным слэшем перед '_', НЕ голым подчёркиванием.""" assert "%/kamensk-uralskiy/%" in _OBLAST_AVITO_URL_PATTERNS assert "%/nizhniy\\_tagil/%" in _OBLAST_AVITO_URL_PATTERNS assert "%/pervouralsk/%" in _OBLAST_AVITO_URL_PATTERNS assert "%/verhnyaya\\_pyshma/%" in _OBLAST_AVITO_URL_PATTERNS assert "%/serov/%" in _OBLAST_AVITO_URL_PATTERNS # ЕКБ обрабатывается отдельным жёстко закодированным паттерном (ekb CTE), # НЕ через этот oblast-список — не должен в него затесаться. assert not any("ekaterinburg" in p for p in _OBLAST_AVITO_URL_PATTERNS) def _like_pattern_to_fnmatch(pattern: str) -> str: """Точный перевод семантики Postgres `LIKE` (default `ESCAPE '\\'`) в fnmatch- паттерн -- посимвольно, а НЕ наивным `.replace()`. LIKE: `%` = любая последовательность символов, `_` = РОВНО один любой символ, `\\%`/`\\_`/`\\\\` = литералы (экранирование). fnmatch: `*` = любая последовательность, `?` = один любой символ; голые `_`/`%` в fnmatch не специальны (можно вставлять как литерал без экранирования). #2578 review: наивный `pat.replace("%", "*")` (как было раньше) НЕ отражал бы семантику `_` вообще -- fnmatch трактует `_` как литерал, LIKE -- как wildcard. Из-за этого расхождения прежний тест не поймал бы латентный баг (нет экранирования `_` в продовых паттернах). Посимвольный разбор здесь корректно различает голый `_` (-> `?` wildcard) и экранированный `\\_` (-> литерал `_`). """ out: list[str] = [] i = 0 n = len(pattern) while i < n: ch = pattern[i] if ch == "\\" and i + 1 < n and pattern[i + 1] in ("%", "_", "\\"): out.append(pattern[i + 1]) # экранированный символ -> литерал as-is i += 2 continue if ch == "%": out.append("*") elif ch == "_": out.append("?") else: out.append(ch) i += 1 return "".join(out) def _in_oblast_or_ekb_scope(source_url: str) -> bool: """Локальная реплика WHERE-условия snapshot-запроса (ekb CTE OR oblast CTE) через корректную LIKE-эмуляцию -- без поднятия БД.""" if fnmatch.fnmatchcase(source_url, _like_pattern_to_fnmatch("%/ekaterinburg/%")): return True return any( fnmatch.fnmatchcase(source_url, _like_pattern_to_fnmatch(pat)) for pat in _OBLAST_AVITO_URL_PATTERNS ) def test_oblast_avito_url_patterns_include_oblast_and_ekb_exclude_foreign_region() -> None: """#2576 DoD: листинг города области и екатеринбургский листинг проходят scope-фильтр; листинг чужого региона (Москва/СПб) — нет. Использует корректную LIKE-эмуляцию (_like_pattern_to_fnmatch), а не наивный `%` -> `*` replace (#2578 review — тот не различал бы `_`-семантику).""" # Область (Каменск-Уральский, #2576 — реальный кейс из тикета) -- проходит. assert _in_oblast_or_ekb_scope("https://www.avito.ru/kamensk-uralskiy/kvartiry/prodam_123") # ЕКБ — по-прежнему проходит (не деградировал). assert _in_oblast_or_ekb_scope("https://www.avito.ru/ekaterinburg/kvartiry/prodam_456") # Чужой регион — НЕ проходит (иначе поехали бы Москва/СПб/Тюмень legacy-строки). assert not _in_oblast_or_ekb_scope("https://www.avito.ru/moskva/kvartiry/prodam_789") assert not _in_oblast_or_ekb_scope("https://www.avito.ru/sankt-peterburg/kvartiry/prodam_000") def test_like_underscore_wildcard_regression_caught_by_escaped_patterns() -> None: """#2578 deep-review latent bug: Postgres `LIKE` трактует `_` как wildcard "ровно один любой символ", а НЕ литерал. Два слага из пяти (nizhniy_tagil, verhnyaya_pyshma) содержат `_` -- БЕЗ экранирования 'nizhniy_tagil' молча совпал бы с 'nizhniyXtagil' (X = любой символ), т.е. коллизия слагов при появлении похожего города. Сегодня коллизий нет (проверено на проде: raw vs escaped паттерны дают одинаковые 776 совпадений), но дыра латентная. Этот тест ДОЛЖЕН падать на RAW (неэкранированном) варианте паттерна -- именно так выглядели продовые паттерны ДО фикса #2578 (`%/nizhniy_tagil/%`, без `\\`). Экранированный прод-паттерн (_OBLAST_AVITO_URL_PATTERNS, ПОСЛЕ фикса) коллизию отклоняет, точный слаг по-прежнему матчит (позитивный кейс жив). """ raw_pattern = "%/nizhniy_tagil/%" # как было бы БЕЗ фикса #2578 (голый '_') escaped_pattern = next(p for p in _OBLAST_AVITO_URL_PATTERNS if "nizhniy" in p) # Сам факт экранирования: прод-паттерн ДОЛЖЕН отличаться от raw ('_' -> '\_'). assert escaped_pattern != raw_pattern, "фикс #2578 должен экранировать '_' в avito_slug" collision_url = "https://www.avito.ru/nizhniyXtagil/kvartiry/prodam_1" exact_url = "https://www.avito.ru/nizhniy_tagil/kvartiry/prodam_1" # RAW: '_' -- wildcard -> ложное совпадение с ЛЮБЫМ символом на его месте. assert fnmatch.fnmatchcase(collision_url, _like_pattern_to_fnmatch(raw_pattern)) # Экранированный прод-паттерн (после фикса): '_' -- литерал -> коллизия отклонена. assert not fnmatch.fnmatchcase(collision_url, _like_pattern_to_fnmatch(escaped_pattern)) # Позитивный кейс не сломан: точный слаг матчит ОБА варианта паттерна. assert fnmatch.fnmatchcase(exact_url, _like_pattern_to_fnmatch(raw_pattern)) assert fnmatch.fnmatchcase(exact_url, _like_pattern_to_fnmatch(escaped_pattern)) @pytest.mark.asyncio async def test_backfill_fetch_exception_continues() -> None: """RuntimeError on one listing -> failed++, db.rollback(), loop continues for next.""" snapshot = _make_snapshot(2) db = _mock_db(snapshot) runs = MagicMock() mock_enrichment = MagicMock() mock_fetch = AsyncMock(side_effect=[RuntimeError("parse error"), mock_enrichment]) fake_settings = _fake_settings( scraper_fetch_mode="cffi", ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=7, params={"batch_size": 10, "budget_sec": 3600} ) assert result.failed == 1 assert result.enriched == 1 assert result.attempted == 2 db.rollback.assert_called() runs.mark_backfill_finished.assert_called_once() @pytest.mark.asyncio async def test_backfill_fetch_timeout_skips_and_continues() -> None: """#1950: один fetch_detail зависает дольше hard-timeout → этот листинг failed, loop НЕ зависает, переходит к следующему, run завершается mark_done (не zombie). Регрессия run 423 (завис 7.7ч → reaped): fetch_detail без timeout блокировал loop навсегда. asyncio.wait_for(timeout) отменяет зависший fetch → TimeoutError. """ snapshot = _make_snapshot(2) db = _mock_db(snapshot) runs = MagicMock() mock_enrichment = MagicMock() call_urls: list[str] = [] async def _fetch( url: str, *, cffi_session: object = None, browser_fetcher: object = None, referer: object = None, reconnect_on_block: bool = True, config: object = None, # Стаб терпит рост сигнатуры fetch_detail: тест про hard-timeout, а не про # набор аргументов. Без этого любой новый kwarg (напр. origin/browser_referer # из #3251) валит его TypeError'ом, хотя к таймауту отношения не имеет. **_extra: object, ): call_urls.append(url) if len(call_urls) == 1: await asyncio.Event().wait() # висит вечно → wait_for отменит по timeout return mock_enrichment # Короткий timeout (50ms) — тест не ждёт реальные 90s; rotate-settle тоже 0. fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_fetch_timeout_s=0.05, avito_proxy_rotate_settle_s=0.0, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, _fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=11, params={"batch_size": 10, "budget_sec": 3600} ) # Первый листинг — failed (timeout), второй — enriched. Run завершён mark_done. assert result.failed == 1 assert result.enriched == 1 assert result.attempted == 2 assert len(call_urls) == 2, "loop должен дойти до второго листинга, а не зависнуть" db.rollback.assert_called() runs.mark_backfill_finished.assert_called_once() runs.mark_failed.assert_not_called() @pytest.mark.asyncio async def test_backfill_listing_gone_marks_inactive_no_breaker() -> None: """#2034: AvitoListingGoneError (мёртвый 404-листинг) → counters.gone++, blocked НЕ растёт, consecutive-block breaker НЕ абортит, листинг помечается is_active=FALSE. Регрессия run 458 (attempted=5 enriched=0 blocked=5 → abort): lat-null очередь состоит из dead-листингов; раньше 404 ловился как soft-block → breaker абортил run до live-листингов. Теперь 404 нейтрален к breaker'у и метит листинг inactive. """ from scraper_kit.avito_exceptions import AvitoListingGoneError snapshot = _make_snapshot(1) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(side_effect=AvitoListingGoneError("404 gone")) mock_scraper = MagicMock() mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True) # use_curl=False: legacy block→_rotate_ip путь (#1551 warm-batch rebuild только при # use_curl=True). MagicMock без явного флага сделал бы use_curl truthy. fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER, mock_scraper), patch(_RUNS, runs), patch(_RESOLVE_PROXY_URL, MagicMock(return_value="http://test-proxy.local:8080")), patch(_FETCH, mock_fetch), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=12, params={"batch_size": 10, "budget_sec": 3600, "max_consecutive_blocks": 5}, ) assert result.gone == 1 assert result.blocked == 0 assert result.attempted == 1 assert result.enriched == 0 # breaker НЕ абортил: rotate_ip НЕ дёргался (gone ≠ block), run завершён mark_done. mock_scraper.return_value._rotate_ip.assert_not_called() runs.mark_backfill_finished.assert_called_once() runs.mark_failed.assert_not_called() # UPDATE listings SET is_active = FALSE по row id=1 выполнен (мок db.execute). update_calls = [ c for c in db.execute.call_args_list if "UPDATE listings SET is_active = FALSE" in str(c.args[0]) ] assert len(update_calls) == 1 assert update_calls[0].args[1] == {"id": 1} _BROWSER_FETCHER = "app.tasks.avito_detail_backfill.BrowserFetcher" @pytest.mark.asyncio async def test_backfill_use_curl_flag_skips_browser_fetcher() -> None: """avito_detail_backfill_use_curl=True: BrowserFetcher не создаётся, fetch_detail вызывается с browser_fetcher=None (curl/backconnect путь). """ snapshot = _make_snapshot(1) db = _mock_db(snapshot) runs = MagicMock() mock_enrichment = MagicMock() mock_fetch = AsyncMock(return_value=mock_enrichment) mock_save = MagicMock(return_value=True) # Флаг use_curl=True, fetch_mode=browser (но флаг перекрывает) fake_settings = _fake_settings( scraper_fetch_mode="browser", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION), patch(_SCRAPER), patch(_BROWSER_FETCHER) as mock_bf_cls, patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, mock_save), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=9, params={"batch_size": 10, "budget_sec": 3600} ) # BrowserFetcher не должен быть создан mock_bf_cls.assert_not_called() # fetch_detail вызван с browser_fetcher=None (curl путь) assert mock_fetch.call_count == 1 _, kwargs = mock_fetch.call_args assert kwargs.get("browser_fetcher") is None assert result.enriched == 1 runs.mark_backfill_finished.assert_called_once() @pytest.mark.asyncio async def test_backfill_use_curl_false_creates_browser_fetcher() -> None: """avito_detail_backfill_use_curl=False + scraper_fetch_mode='browser': BrowserFetcher создаётся (legacy browser/auv путь). """ snapshot = _make_snapshot(1) db = _mock_db(snapshot) runs = MagicMock() mock_enrichment = MagicMock() mock_fetch = AsyncMock(return_value=mock_enrichment) mock_save = MagicMock(return_value=True) fake_settings = _fake_settings( scraper_fetch_mode="browser", avito_detail_backfill_use_curl=False, ) mock_bf_instance = AsyncMock() mock_bf_instance.__aenter__ = AsyncMock(return_value=mock_bf_instance) mock_bf_instance.__aexit__ = AsyncMock(return_value=False) with ( patch(_SETTINGS, fake_settings), patch(_SESSION), patch(_SCRAPER), patch(_BROWSER_FETCHER, return_value=mock_bf_instance) as mock_bf_cls, patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, mock_save), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=10, params={"batch_size": 10, "budget_sec": 3600} ) # BrowserFetcher должен быть создан (source="avito", endpoint из settings — #2310 # kit BrowserFetcher требует endpoint= обязательным keyword-only параметром). # # Проверяем ВХОЖДЕНИЕ аргументов, а не полное равенство вызова. Прежняя # форма `assert_called_once_with(source=..., endpoint=...)` требовала, чтобы # других аргументов НЕ БЫЛО, и тем самым закрепляла дефект: именно # отсутствие проводки пула (#3045) уводило браузерный прогон мимо прокси — # 5 блоков из 5 попыток при полностью здоровом пуле. Проводка проверяется # отдельно ниже. mock_bf_cls.assert_called_once() _bf_kwargs = mock_bf_cls.call_args.kwargs assert _bf_kwargs["source"] == "avito" assert _bf_kwargs["endpoint"] == fake_settings.browser_http_endpoint # #3045: без этих трёх фетчер не кладёт "proxy" в тело POST /fetch, и # сайдкар уходит на свой env-прокси мимо пула (proxy_lease_id=None). assert _bf_kwargs.get("proxy_provider") is not None, "пул не подключён" assert "use_pool" in _bf_kwargs, "флаг пула не доезжает до фетчера" assert "environment" in _bf_kwargs, "без него отказ «пул пуст» мёртв (#2616)" # fetch_detail вызван с browser_fetcher установленным (не None) assert mock_fetch.call_count == 1 _, kwargs = mock_fetch.call_args assert kwargs.get("browser_fetcher") is not None assert result.enriched == 1 runs.mark_backfill_finished.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 scraper_kit.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 = _fake_settings( 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_backfill_finished.assert_called_once() runs.mark_failed.assert_not_called() # --------------------------------------------------------------------------- # Отказы-не-блоки: брейкер + перепись причин (2026-08-06) # --------------------------------------------------------------------------- @pytest.mark.asyncio async def test_backfill_aborts_on_consecutive_failures_and_names_the_reason() -> None: """Серия отказов без единого успеха обрывается, а причина попадает в статус прогона. Прод 3-5 августа: attempted≈1600, failed≈1600, blocked=0, enriched=0, весь бюджет 9000 с и 1600 запросов через единственный прокси — и ни слова о том, ЧТО именно отказало (поштучные отказы логируются WARNING, а логи контейнера пропадают на первом деплое). Брейкера на отказы-не-блоки не было вовсе. """ snapshot = _make_snapshot(200) db = _mock_db(snapshot) runs = MagicMock() # Тот же класс отказа, что видели у соседнего свипа в тот же день. mock_fetch = AsyncMock( side_effect=OSError( "Failed to perform, curl: (56) CONNECT tunnel failed, response 502. " "See https://curl.se/libcurl/c/libcurl-errors.html" ) ) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=77, params={ "batch_size": 200, "budget_sec": 3600, "max_consecutive_failures": 25, }, ) assert result.attempted == 25, "серия отказов обязана обрываться, а не выедать бюджет" assert result.failed == 25 assert result.blocked == 0 hint = runs.mark_backfill_finished.call_args.kwargs["fail_hint"] assert hint is not None assert "OSError" in hint # тип исключения = кому принадлежит отказ assert "CONNECT tunnel failed" in hint assert "25 из 25" in hint # доля, а не единичный пример assert "https://curl.se" not in hint # URL вырезан, иначе 1600 «разных» причин @pytest.mark.asyncio async def test_backfill_success_resets_failure_streak() -> None: """Успех между отказами обнуляет серию — здоровый прогон брейкер не трогает.""" snapshot = _make_snapshot(5) db = _mock_db(snapshot) runs = MagicMock() boom = ValueError("avito detail HTTP 500 for https://www.avito.ru/x") # 2 отказа, успех, 2 отказа — при пороге 3 ни одна серия его не достигает. mock_fetch = AsyncMock(side_effect=[boom, boom, MagicMock(), boom, boom]) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=78, params={"batch_size": 5, "budget_sec": 3600, "max_consecutive_failures": 3}, ) assert result.attempted == 5 assert result.enriched == 1 assert result.failed == 4 def _spaced_outcomes(total: int, block_positions: set[int]) -> list: """1-indexed позиции block_positions -> AvitoBlockedError, остальные -- success MagicMock. Ровный интервал между блоками (не пачками) -- моделирует прод-профиль 4786/5200 (#3184): блоки есть, но локально в любом 20-окне их доля намного ниже порога. """ return [ AvitoBlockedError("ip blocked") if pos in block_positions else MagicMock() for pos in range(1, total + 1) ] @pytest.mark.asyncio async def test_backfill_25pct_blocks_spread_does_not_abort() -> None: """47 блоков из 190 попыток (25%, прогон 4786) -- НЕ обрывается (#3184). Раньше это же соотношение (после единственной пачки >= max_consecutive_blocks) ловилось общим брейкером и уходило в 'banned', хотя прогон честно обогатил 125/190 карточек. Блоки расставлены каждый 4-й (изолированно, не пачкой) -- ни в одном 20-окне доля не приближается к порогу 0.7. """ total = 190 block_positions = {pos for pos in range(1, total + 1) if pos % 4 == 0} assert len(block_positions) == 47 snapshot = _make_snapshot(total) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions)) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=101, params={"batch_size": total, "budget_sec": 3600} ) assert result.attempted == total assert result.blocked == 47 assert result.enriched == total - 47 runs.mark_backfill_finished.assert_called_once() assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is False @pytest.mark.asyncio async def test_backfill_32pct_blocks_spread_does_not_abort() -> None: """20 блоков из 63 попыток (32%, прогон 5200) -- НЕ обрывается (#3184).""" total = 63 block_positions = {pos for pos in range(3, 61, 3)} assert len(block_positions) == 20 snapshot = _make_snapshot(total) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions)) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=102, params={"batch_size": total, "budget_sec": 3600} ) assert result.attempted == total assert result.blocked == 20 assert result.enriched == total - 20 runs.mark_backfill_finished.assert_called_once() assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is False @pytest.mark.asyncio async def test_backfill_burst_mid_healthy_run_does_not_abort_and_streak_logged() -> None: """6 блоков подряд посреди здорового прогона -- НЕ обрывает (#3184), а пачка едет в counters["block_streak_histogram"]. Раньше max_consecutive_blocks=5 абортил бы уже на 5-м блоке пачки независимо от контекста; серия сама по себе больше не критерий -- 20 успехов до пачки держат окно далеко от порога 0.7. """ total = 50 block_positions = set(range(21, 27)) # 6 подряд, позиции 21..26 assert len(block_positions) == 6 snapshot = _make_snapshot(total) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions)) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=103, params={"batch_size": total, "budget_sec": 3600} ) assert result.attempted == total assert result.blocked == 6 runs.mark_backfill_finished.assert_called_once() assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is False counters_arg = runs.mark_backfill_finished.call_args.args[2] assert counters_arg["block_streak_histogram"] == {"6": 1} # jsonb -- ключи строкой @pytest.mark.asyncio async def test_backfill_ratio_criterion_aborts_after_window_fills_with_blocks() -> None: """Успех на 1-й попытке, дальше сплошные блоки -- обрывает ПО ДОЛЕ в окне, не safety-net (#3184 review MAJOR 3: раньше ни один тест не проверял срабатывание ratio-критерия, только его отсутствие). Окно (20) заполняется целиком ровно на 20-й попытке (1 успех + 19 блоков) -- ratio=19/20=0.95 >= порога 0.7, абортит там же, до 21-й попытки дело не доходит. Safety-net здесь не может сработать: success на 1-й попытке снимает _pure_block_run сразу, до любого блока. """ total = 200 block_positions = set(range(2, total + 1)) # всё, кроме 1-й попытки snapshot = _make_snapshot(total) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions)) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=104, params={"batch_size": total, "budget_sec": 3600} ) assert result.attempted == 20 assert result.blocked == 19 assert result.enriched == 1 runs.mark_backfill_finished.assert_called_once() assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is True @pytest.mark.asyncio async def test_backfill_ratio_boundary_14_of_20_aborts() -> None: """14/20 = 0.7 -- ровно порог, should_abort сравнивает через >=, значит абортит (#3184 review MAJOR 3 -- граница снизу).""" total = 20 block_positions = set(range(7, 21)) # позиции 7..20 включительно -- 14 штук assert len(block_positions) == 14 snapshot = _make_snapshot(total) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions)) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=105, params={"batch_size": total, "budget_sec": 3600} ) assert result.blocked == 14 runs.mark_backfill_finished.assert_called_once() assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is True @pytest.mark.asyncio async def test_backfill_ratio_boundary_13_of_20_does_not_abort() -> None: """13/20 = 0.65 -- ниже порога 0.7, НЕ абортит (#3184 review MAJOR 3 -- граница сверху, зеркало предыдущего теста).""" total = 20 block_positions = set(range(8, 21)) # позиции 8..20 включительно -- 13 штук assert len(block_positions) == 13 snapshot = _make_snapshot(total) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions)) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), ): result = await run_avito_detail_backfill( db, run_id=106, params={"batch_size": total, "budget_sec": 3600} ) assert result.attempted == 20 assert result.blocked == 13 runs.mark_backfill_finished.assert_called_once() assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is False @pytest.mark.asyncio async def test_backfill_ratio_abort_log_names_the_ratio_not_the_streak(caplog: Any) -> None: """ABORT по доле называет долю, а не текущую серию (#3184 follow-up). Воспроизводит прод-прогон 5210: 14 блоков из 20, обрыв ровно по порогу 0.7, но ПОСЛЕДНЯЯ серия блоков к этому моменту равна единице (позиция 19 -- успех, позиция 20 -- блок). Лог печатал "ABORT -- 1 consecutive blocks": число верное, величина не та. Читатель видит цифру, по которой обрыва быть не могло. """ total = 20 # Раскладка ровно как в прогоне 5210. block_positions = {2, 3, 5, 6, 7, 8, 9, 11, 12, 13, 14, 16, 17, 20} assert len(block_positions) == 14 assert 19 not in block_positions and 20 in block_positions, "серия на обрыве = 1" snapshot = _make_snapshot(total) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions)) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), caplog.at_level("ERROR"), ): result = await run_avito_detail_backfill( db, run_id=110, params={"batch_size": total, "budget_sec": 3600} ) assert result.blocked == 14 abort_records = [r.message for r in caplog.records if "ABORT" in r.message] assert abort_records, "ожидался ABORT-лог" msg = abort_records[0] assert "14/20" in msg, msg assert "consecutive" not in msg, msg assert "1 блоков подряд" not in msg, msg counters = runs.mark_backfill_finished.call_args.args[2] assert counters["abort_reason"] == "ratio" @pytest.mark.asyncio async def test_backfill_safety_net_abort_log_names_the_streak(caplog: Any) -> None: """ABORT по safety-net называет серию -- ту величину, по которой он и сработал. Зеркало предыдущего теста: снапшот короче окна, ratio недостижим, рвёт серия. Если бы обе ветки печатали один и тот же текст, тест выше проходил бы и у сломанного сообщения. """ total = 5 # короче окна 20 -- ratio-критерий физически недостижим snapshot = _make_snapshot(total) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(side_effect=AvitoBlockedError("ip blocked")) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), caplog.at_level("ERROR"), ): result = await run_avito_detail_backfill( db, run_id=111, params={"batch_size": total, "budget_sec": 3600} ) assert result.blocked == 5 abort_records = [r.message for r in caplog.records if "ABORT" in r.message] assert abort_records, "ожидался ABORT-лог" msg = abort_records[0] assert "5 блоков подряд" in msg, msg assert "доля блоков" not in msg, msg counters = runs.mark_backfill_finished.call_args.args[2] assert counters["abort_reason"] == "safety_net" @pytest.mark.asyncio async def test_backfill_no_abort_leaves_no_abort_reason_in_counters() -> None: """Прогон без обрыва не пишет abort_reason -- ключ означает 'оборвались', а не 'считали критерий'.""" total = 20 block_positions = {8, 9, 10, 11, 12, 13, 14, 15, 16, 17, 18, 19, 20} # 13/20 = 0.65 assert len(block_positions) == 13 snapshot = _make_snapshot(total) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions)) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION, return_value=AsyncMock()), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), ): await run_avito_detail_backfill( db, run_id=112, params={"batch_size": total, "budget_sec": 3600} ) counters = runs.mark_backfill_finished.call_args.args[2] assert "abort_reason" not in counters # --------------------------------------------------------------------------- # Rotate-by-attempt-counter (задача поверх #3298) # --------------------------------------------------------------------------- def _browser_mock_instance(lease_id: int | None) -> AsyncMock: """Мокнутый BrowserFetcher instance с lease_id и sync request_context_reset (реальный метод -- НЕ async, mock должен звать его так же).""" inst = AsyncMock() inst.__aenter__ = AsyncMock(return_value=inst) inst.__aexit__ = AsyncMock(return_value=False) inst.lease_id = lease_id inst.request_context_reset = MagicMock() return inst @pytest.mark.asyncio async def test_rotate_triggers_at_threshold() -> None: """attempts_since_rotation достигает порога -> rotate_proxy вызван РОВНО один раз с proxy_id текущего lease (browser_fetcher.lease_id).""" snapshot = _make_snapshot(3) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(return_value=MagicMock()) mock_bf_instance = _browser_mock_instance(lease_id=777) fake_settings = _fake_settings( scraper_fetch_mode="browser", avito_detail_backfill_use_curl=False, avito_detail_backfill_rotate_after_attempts=3, ) mock_rotate = AsyncMock( return_value=RotationResult(ok=True, reason=None, new_ip="1.2.3.4", reconnect_delay_s=4.0) ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION), patch(_SCRAPER), patch(_BROWSER_FETCHER, return_value=mock_bf_instance), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), patch(_ROTATE_PROXY, mock_rotate), ): await run_avito_detail_backfill( db, run_id=200, params={"batch_size": 3, "budget_sec": 3600} ) mock_rotate.assert_called_once() args, _ = mock_rotate.call_args assert args[1] == 777 # proxy_id из lease_id @pytest.mark.asyncio async def test_rotate_resets_counter_after_trigger() -> None: """Порог=2, снапшот=4 -> ротация срабатывает ДВАЖДЫ (счётчик обнуляется после каждого срабатывания, не только один раз за прогон).""" snapshot = _make_snapshot(4) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(return_value=MagicMock()) mock_bf_instance = _browser_mock_instance(lease_id=5) fake_settings = _fake_settings( scraper_fetch_mode="browser", avito_detail_backfill_use_curl=False, avito_detail_backfill_rotate_after_attempts=2, ) mock_rotate = AsyncMock( return_value=RotationResult(ok=True, reason=None, new_ip="9.9.9.9", reconnect_delay_s=1.0) ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION), patch(_SCRAPER), patch(_BROWSER_FETCHER, return_value=mock_bf_instance), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), patch(_ROTATE_PROXY, mock_rotate), ): await run_avito_detail_backfill( db, run_id=201, params={"batch_size": 4, "budget_sec": 3600} ) assert mock_rotate.call_count == 2 @pytest.mark.asyncio async def test_rotate_waits_reconnect_delay_from_provider() -> None: """reconnect_delay_s из RotationResult -- ждём именно его, не дефолт.""" snapshot = _make_snapshot(1) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(return_value=MagicMock()) mock_bf_instance = _browser_mock_instance(lease_id=1) fake_settings = _fake_settings( scraper_fetch_mode="browser", avito_detail_backfill_use_curl=False, avito_detail_backfill_rotate_after_attempts=1, ) mock_rotate = AsyncMock( return_value=RotationResult(ok=True, reason=None, new_ip="1.1.1.1", reconnect_delay_s=6.5) ) mock_sleep = AsyncMock() with ( patch(_SETTINGS, fake_settings), patch(_SESSION), patch(_SCRAPER), patch(_BROWSER_FETCHER, return_value=mock_bf_instance), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, mock_sleep), patch(_ROTATE_PROXY, mock_rotate), ): await run_avito_detail_backfill( db, run_id=202, params={"batch_size": 1, "budget_sec": 3600} ) sleep_values = [call.args[0] for call in mock_sleep.call_args_list] assert 6.5 in sleep_values @pytest.mark.asyncio async def test_rotate_uses_default_delay_when_provider_silent() -> None: """reconnect_delay_s=None (провайдер не прислал rt) -> используем дефолт _DEFAULT_ROTATE_RECONNECT_DELAY_S, не падаем и не ждём 0с.""" snapshot = _make_snapshot(1) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(return_value=MagicMock()) mock_bf_instance = _browser_mock_instance(lease_id=1) fake_settings = _fake_settings( scraper_fetch_mode="browser", avito_detail_backfill_use_curl=False, avito_detail_backfill_rotate_after_attempts=1, ) mock_rotate = AsyncMock( return_value=RotationResult(ok=True, reason=None, new_ip="2.2.2.2", reconnect_delay_s=None) ) mock_sleep = AsyncMock() with ( patch(_SETTINGS, fake_settings), patch(_SESSION), patch(_SCRAPER), patch(_BROWSER_FETCHER, return_value=mock_bf_instance), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, mock_sleep), patch(_ROTATE_PROXY, mock_rotate), ): await run_avito_detail_backfill( db, run_id=203, params={"batch_size": 1, "budget_sec": 3600} ) sleep_values = [call.args[0] for call in mock_sleep.call_args_list] assert _DEFAULT_ROTATE_RECONNECT_DELAY_S in sleep_values @pytest.mark.asyncio async def test_rotate_success_calls_context_reset() -> None: """Успешная ротация -> request_context_reset() вызван (новый адрес + свежий browser-контекст, старые куки выброшены).""" snapshot = _make_snapshot(1) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(return_value=MagicMock()) mock_bf_instance = _browser_mock_instance(lease_id=42) fake_settings = _fake_settings( scraper_fetch_mode="browser", avito_detail_backfill_use_curl=False, avito_detail_backfill_rotate_after_attempts=1, ) mock_rotate = AsyncMock( return_value=RotationResult(ok=True, reason=None, new_ip="3.3.3.3", reconnect_delay_s=1.0) ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION), patch(_SCRAPER), patch(_BROWSER_FETCHER, return_value=mock_bf_instance), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), patch(_ROTATE_PROXY, mock_rotate), ): await run_avito_detail_backfill( db, run_id=204, params={"batch_size": 1, "budget_sec": 3600} ) mock_bf_instance.request_context_reset.assert_called_once() @pytest.mark.asyncio async def test_rotate_failure_does_not_abort_run_and_skips_context_reset() -> None: """rotate_proxy(ok=False) (лимит исчерпан / нет rotate_url / провайдер отказал) -- прогон НЕ падает, продолжает на текущем адресе, request_context_reset НЕ вызывается (менять контекст без нового IP бессмысленно).""" snapshot = _make_snapshot(2) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(return_value=MagicMock()) mock_save = MagicMock(return_value=True) mock_bf_instance = _browser_mock_instance(lease_id=9) fake_settings = _fake_settings( scraper_fetch_mode="browser", avito_detail_backfill_use_curl=False, avito_detail_backfill_rotate_after_attempts=1, ) mock_rotate = AsyncMock( return_value=RotationResult(ok=False, reason="daily rotation limit reached (200/day)") ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION), patch(_SCRAPER), patch(_BROWSER_FETCHER, return_value=mock_bf_instance), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, mock_save), patch(_SLEEP, new_callable=AsyncMock), patch(_ROTATE_PROXY, mock_rotate), ): result = await run_avito_detail_backfill( db, run_id=205, params={"batch_size": 2, "budget_sec": 3600} ) # прогон не упал и довёл оба листинга до конца assert result.enriched == 2 assert mock_rotate.call_count == 2 # счётчик обнулился и после неудачи -> 2 попытки mock_bf_instance.request_context_reset.assert_not_called() runs.mark_backfill_finished.assert_called_once() # rotate НЕ считается блоком площадки assert result.blocked == 0 @pytest.mark.asyncio async def test_rotate_not_counted_as_platform_block() -> None: """Ротация не трогает breaker/blocked/ban_kinds -- это обслуживание канала, не исход fetch(). aborted_by_blocks остаётся False, blocked=0 даже когда rotate срабатывает на каждой попытке.""" snapshot = _make_snapshot(5) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(return_value=MagicMock()) mock_bf_instance = _browser_mock_instance(lease_id=3) fake_settings = _fake_settings( scraper_fetch_mode="browser", avito_detail_backfill_use_curl=False, avito_detail_backfill_rotate_after_attempts=1, ) mock_rotate = AsyncMock( return_value=RotationResult(ok=True, reason=None, new_ip="4.4.4.4", reconnect_delay_s=0.1) ) with ( patch(_SETTINGS, fake_settings), patch(_SESSION), patch(_SCRAPER), patch(_BROWSER_FETCHER, return_value=mock_bf_instance), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), patch(_ROTATE_PROXY, mock_rotate), ): result = await run_avito_detail_backfill( db, run_id=206, params={"batch_size": 5, "budget_sec": 3600} ) assert mock_rotate.call_count == 5 assert result.blocked == 0 finished_kwargs = runs.mark_backfill_finished.call_args.kwargs assert finished_kwargs.get("aborted_by_blocks") is False assert not finished_kwargs.get("ban_kinds") @pytest.mark.asyncio async def test_rotate_skipped_in_curl_mode_no_lease() -> None: """use_curl=True (прод-дефолт): browser_fetcher остаётся None -- rotate_proxy НЕ вызывается вообще (curl/backconnect путь не держит lease с proxy_id).""" snapshot = _make_snapshot(20) db = _mock_db(snapshot) runs = MagicMock() mock_fetch = AsyncMock(return_value=MagicMock()) fake_settings = _fake_settings( scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True, avito_detail_backfill_rotate_after_attempts=1, ) mock_rotate = AsyncMock() with ( patch(_SETTINGS, fake_settings), patch(_SESSION), patch(_SCRAPER), patch(_RUNS, runs), patch(_FETCH, mock_fetch), patch(_SAVE, return_value=True), patch(_SLEEP, new_callable=AsyncMock), patch(_ROTATE_PROXY, mock_rotate), ): await run_avito_detail_backfill( db, run_id=207, params={"batch_size": 20, "budget_sec": 3600} ) mock_rotate.assert_not_called()