"""Scheduled backfill: detail-enrichment for legacy avito listings (#1551). Nightly window 09:00-12:00 UTC (migration 112, source=avito_detail_backfill). Offset from avito_city_sweep (~03-04 UTC) to avoid sharing proxy IP simultaneously. Problem: ~15243/15917 avito listings have detail_enriched_at IS NULL. City sweep ENRICH_DETAIL only enriches top-N proximate listings per anchor. Legacy listings (older than 2h or outside radius) are never enriched. Solution: single snapshot SELECT at start (guarantees termination), same proxy session path as the detail-phase of `run_avito_city_sweep` (scraper_kit.orchestration.pipeline). Block handling mirrors that phase: rotate IP on every block, abort after max_consecutive_blocks. Статус оборванного блоками прогона — 'banned' (#2674, runs.mark_backfill_finished): работу он не доделал, остаток снапшота уедет в следующую ночь через NULL detail_enriched_at. Отказы, не являющиеся блоками, до 2026-08-06 брейкера не имели вовсе: прогоны 3-5 августа делали ~1600 попыток, получали 1600 отказов, ноль обогащений и выедали весь бюджет (9000 с) вместе с 1600 запросами через единственный прокси. Теперь такая серия обрывается по max_consecutive_failures, а самая частая причина отказа пишется в текст статуса прогона (_failure_signature) — иначе она живёт только в логах контейнера, а те исчезают на первом же деплое. ДИАГНОЗ бана (#2764): каждый пойманный блок классифицируется по ТИПУ исключения (ban_kind_of_exception) и уходит в scrape_runs.ban_kind. До этой правки прогон 3306 (blocked=5, 0 обогащено) получил 'platform' по УМОЛЧАНИЮ — финализатор диагноз не передавал, а browser-режим fetch_detail всё равно превращал отказ сайдкара в AvitoBlockedError, так что передавать было бы нечего. КРИТЕРИЙ ОБРЫВА (#3184): голое "N блоков подряд" (max_consecutive_blocks) абортило прогон 125/190 (25% блоков, честно обогащённых 125 карточек) статусом 'banned' наравне с прогоном 0/5 (100% блоков) — блоки идут автокоррелированными пачками, и пачка 5-6 подряд у здорового прогона с базовой долей 25-46% не редкость. Теперь абортит app.services.backfill_block_breaker.BlockRatioBreaker: доля блоков в скользящем окне последних N попыток (settings.detail_backfill_ block_ratio_window/_threshold) ИЛИ safety-net для коротких прогонов (max_ consecutive_blocks блоков подряд БЕЗ единого успеха с самого начала — старое поведение для случая вида 0/5). Гистограмма длин пачек блоков едет в counters["block_streak_histogram"], иначе эффект правки нечем измерить постфактум. Какой из двух критериев сработал -- в counters["abort_reason"] ("ratio" / "safety_net"); ключа нет, если прогон не обрывался. Он же называется в ABORT-логе своей величиной: доля печатает "14/20", safety-net -- длину серии. Раньше лог печатал серию всегда, и прогон 5210 (обрыв по доле 14/20) отчитался как "ABORT -- 1 consecutive blocks". """ from __future__ import annotations import asyncio import logging import random import re import time from collections import Counter from dataclasses import dataclass, field from urllib.parse import urlparse from curl_cffi.requests import AsyncSession from scraper_kit.avito_exceptions import ( AvitoBlockedError, AvitoListingGoneError, AvitoRateLimitedError, ) from scraper_kit.browser_fetcher import BrowserFetcher from scraper_kit.orchestration.pipeline import CITY_LOCATIONS, ban_kind_of_exception # #2397 slice B (эпик #2277 decommission scrape_pipeline.py, Part E): раньше # _CHROME_HEADERS/_avito_proxies() импортировались из app.services.scrape_pipeline. # kit уже вынес byte-идентичный dict (DOCUMENT_HEADERS) и ту же формулу # {"http": url, "https": url} if url else None (http_proxies) в # providers/_base.py (#2358 Foundation) — реюзаем вместо копии, разрывая # зависимость от scrape_pipeline.py перед его удалением. from scraper_kit.providers._base import DEFAULT_IMPERSONATE, DOCUMENT_HEADERS, http_proxies from scraper_kit.providers.avito.detail import ( _AVITO_WARM_SEARCH_URL, _serp_origin_for, build_warmed_session, fetch_detail, research_in_session, save_detail_enrichment, ) from scraper_kit.providers.avito.serp import AvitoScraper from scraper_kit.proxy_errors import NoProxyAvailableError from scraper_kit.snapshot_writer import upsert_listing_snapshot from sqlalchemy import text from sqlalchemy.orm import Session from app.core.config import settings from app.core.shutdown import shutdown_requested from app.services import scrape_runs as runs_mod from app.services.backfill_block_breaker import BlockRatioBreaker from app.services.proxy_egress import resolve_proxy_url from app.services.proxy_rotation import rotate_proxy from app.services.scraper_adapters import RealProxyProvider, RealScraperConfig # #2397 Part D1 (#2330 закрыт): _AVITO_WARM_SEARCH_URL/build_warmed_session больше # НЕ legacy — kit's scraper_kit.providers.avito.detail.build_warmed_session() теперь # принимает config: ScraperConfig | None (Strangler-инъекция #2330) и пробрасывает # его в _build_detail_session() для settings.scraper_proxy_url (sticky МГТС-прокси), # зеркаля fetch_detail. Оба call site'а ниже передают config=RealScraperConfig() — # без него warm-batch curl-путь (use_curl=True, ДЕФОЛТ прод-режима, # avito_detail_backfill_use_curl=True) молча терял бы backconnect (прямое # datacenter-подключение вместо sticky-прокси), тот же класс бага что #2322/#2310. logger = logging.getLogger(__name__) __all__ = [ "AvitoDetailBackfillResult", "run_avito_detail_backfill", ] # #2576 этап B: oblast-города (region 66, вне ЕКБ) уже дают листинги (Каменск- # Уральский), но snapshot-SELECT ниже раньше фильтровал ЖЁСТКО '%/ekaterinburg/%' — # у всех остальных detail_enriched_at оставался NULL навсегда (без detail-страницы # нет lat/lon -> листинг молча выпадает из подбора аналогов по радиусу). # CITY_LOCATIONS.avito_slug — единственный источник правды для avito URL-слага # города (может отличаться от нашего city_slug: kamensk-uralskiy через дефис, # verhnyaya_pyshma без "kh") -- дублировать список тут вместо импорта было бы # risk дрейфа при добавлении новых oblast-городов. # # #2578 review: Postgres LIKE трактует '_' как wildcard "один любой символ" (не # литерал) и '%' как wildcard "любая последовательность" -- два слага из пяти # (nizhniy_tagil, verhnyaya_pyshma) содержат '_', без экранирования это латентная # дыра: город с похожим слагом (напр. nizhniyXtagil) молча совпал бы. Сегодня # коллизий нет (проверено на проде: raw vs escaped паттерны дают одинаковые 776 # совпадений), но экранируем сейчас, а не когда появится реальная коллизия. # LIKE по умолчанию использует '\' как escape-символ (без явного ESCAPE) — # подтверждено на живом Postgres 16.4 (см. коммит #2578-fixup): 'nizhniyXtagil' # матчит неэкранированный '%/nizhniy_tagil/%' (LIKE default '_'=wildcard) и НЕ # матчит экранированный '%/nizhniy\_tagil/%' (LIKE '\_' = литерал '_'); точный # слаг 'nizhniy_tagil' матчит оба варианта -- позитивный кейс не сломан. # # #262 wave 2: avito_slug — Optional в CityLocation (не у каждого oblast-города # подтверждён). Города без avito_slug пропускаем целиком — у них НЕТ avito_city_ # sweep schedule (262_), значит НЕТ и avito-листингов с их URL; паттерн для них # был бы либо мёртвым, либо (что хуже) построен из city_slug вместо реального # avito URL-сегмента и создал бы ложный LIKE-матч. _OBLAST_AVITO_URL_PATTERNS = tuple( "%/" + loc.avito_slug.replace("\\", "\\\\").replace("_", "\\_").replace("%", "\\%") + "/%" for loc in CITY_LOCATIONS.values() if loc.avito_slug is not None ) # Причина отказа карточки без её URL: 1576 отказов одного прогона должны схлопнуться # в ОДНУ строку, иначе перепись бесполезна. _URL_IN_MESSAGE_RE = re.compile(r"https?://\S+") def _failure_signature(exc: BaseException) -> str: """Подпись причины отказа: тип исключения + текст без URL. Зачем (замер 2026-08-06): у прогонов 3 и 4 августа counters говорили `attempted=1576, failed=1576, blocked=0` — и ничего больше. Кто отказал, площадка или наш тракт, было видно ТОЛЬКО в логах контейнера, а тот пересоздаётся на каждом деплое и уносит их с собой; в GlitchTip попадают события уровня ERROR, а поштучные отказы — WARNING. Разница между этими двумя диагнозами — разные владельцы задачи (#2686, #2698), поэтому она обязана переживать перезапуск контейнера, то есть лежать в самом прогоне. Тип исключения — первый разряд диагноза (AvitoBlockedError = площадка показала 403/firewall; сетевой класс curl_cffi = наш прокси-тракт; ValueError = ответ пришёл, но не разобран), текст — второй. """ message = _URL_IN_MESSAGE_RE.sub("", str(exc)).strip() return f"{type(exc).__name__}: {message}"[:160] if message else type(exc).__name__ def _top_failure(census: Counter[str]) -> str | None: """Самая частая причина отказа с её долей; None — отказов не было.""" if not census: return None reason, hits = census.most_common(1)[0] return f"{reason} ({hits} из {sum(census.values())})" def _iter_causes(exc: BaseException) -> list[BaseException]: """Цепочка причин исключения, без зацикливания (копия приёма из domclick_detail_backfill — там же он и обкатан на #3283).""" seen: set[int] = set() out: list[BaseException] = [] cur: BaseException | None = exc while cur is not None and id(cur) not in seen: out.append(cur) seen.add(id(cur)) cur = cur.__cause__ or cur.__context__ return out def _caused_by_empty_pool(exc: BaseException) -> bool: """Прячется ли за этим «блоком» пустой пул прокси (#3288, как #3283 у домклика). `NoProxyAvailableError` документирован ровно как «НАША инфраструктура, не внешний блок», и поднимается ДО HTTP-запроса: к площадке мы не ходили вовсе. Сюда он попадает под видом блокировки, потому что fetch_detail заворачивает в `AvitoSidecarUnavailableError` любое исключение фетча. Опора — ТИП в цепочке `__cause__`/`__context__`, а не подстрока «no proxy available» в тексте: текст обёртки её действительно содержит, но ровно так #3272 уже один раз объявил блоком пользовательское описание квартиры. """ return any(isinstance(c, NoProxyAvailableError) for c in _iter_causes(exc)) # Провайдер mobileproxy.space почти всегда возвращает "rt" (секунды на переподключение # канала после ротации) в RotationResult.reconnect_delay_s -- см. app.services. # proxy_rotation docstring. Дефолт нужен ТОЛЬКО если провайдер его не прислал # (best-effort парсинг "rt", неудача не является ошибкой ротации) -- замеренный # вживую диапазон простоя канала 2-12с, 10с чуть выше нижней границы. _DEFAULT_ROTATE_RECONNECT_DELAY_S = 10.0 async def _rotate_current_proxy( db: Session, run_id: int, browser_fetcher: BrowserFetcher, trigger: str = "rotate-by-attempts", ) -> bool: """Ротация exit-IP арендованного прокси + сброс browser-контекста. Два триггера, оба пишут свой ярлык в логи через ``trigger``: плановый по счётчику попыток (``rotate-by-attempts``) и реактивный по бану площадки (``rotate-on-ban``, #3283). Разделение нужно только для читаемости логов — механика одна и та же. Свежий адрес со старыми куками бесполезен — личность браузера должна меняться ВМЕСТЕ с адресом, иначе следующий запрос уходит с нового IP, но с cookie-следом старого (request_context_reset потребляется ровно следующим fetch(), #3118). Работает ТОЛЬКО когда есть browser_fetcher (browser-режим) — это единственный путь, где avito_detail_backfill реально держит lease с proxy_id (scrape_proxies.id, browser_fetcher.lease_id). curl/backconnect-путь (use_curl=True) идёт через sticky settings.scraper_proxy_url БЕЗ lease — там ротировать нечего, вызывающий цикл не зовёт эту функцию в том режиме вовсе. Прод при этом в browser-режиме, а не в curl: у контейнера tradein-scraper (там же живёт планировщик) проверено AVITO_DETAIL_BACKFILL_USE_CURL=false при SCRAPER_FETCH_MODE=browser и USE_PROXY_POOL_BROWSER=true, так что ротация активна. Значение true стоит только у tradein-backend, который добор не запускает. Комментарий ниже по файлу (~строка 653) называет use_curl=True «прод-дефолтом» — это предсуществующее заблуждение, а не описание текущего прода. Отказ провайдера (лимит исчерпан, нет rotate_url, сетевой сбой, неизвестный хост) НЕ должен ронять прогон — логируем и продолжаем на текущем адресе. Эта попытка НЕ является блоком/отказом площадки и не должна попадать в BlockRatioBreaker/ counters.blocked/ban_kinds — вызывающий код не передаёт её исход ни в один из них, это обслуживание канала, а не результат fetch(). """ proxy_id = browser_fetcher.lease_id if proxy_id is None: logger.info( "avito_detail_backfill: run_id=%d %s skipped -- no leased proxy", run_id, trigger, ) return False result = await rotate_proxy(db, proxy_id) if not result.ok: logger.warning( "avito_detail_backfill: run_id=%d %s FAILED proxy_id=%d: %s", run_id, trigger, proxy_id, result.reason, ) return False delay = ( result.reconnect_delay_s if result.reconnect_delay_s is not None else _DEFAULT_ROTATE_RECONNECT_DELAY_S ) logger.info( "avito_detail_backfill: run_id=%d %s OK proxy_id=%d new_ip=%s -- " "waiting %.1fs for channel reconnect", run_id, trigger, proxy_id, result.new_ip, delay, ) await asyncio.sleep(delay) browser_fetcher.request_context_reset() return True @dataclass class AvitoDetailBackfillResult: """Counters for one backfill run.""" attempted: int = 0 enriched: int = 0 blocked: int = 0 gone: int = 0 failed: int = 0 duration_sec: float = field(default=0.0) def to_dict(self) -> dict[str, int]: return { "attempted": self.attempted, "enriched": self.enriched, "blocked": self.blocked, "gone": self.gone, "failed": self.failed, "duration_sec": int(self.duration_sec), } async def run_avito_detail_backfill( db: Session, *, run_id: int, params: dict, ) -> AvitoDetailBackfillResult: """Backfill detail_enriched_at for legacy avito listings via mobile proxy. Params (from default_params jsonb in scrape_schedules): batch_size: int -- ЕКБ snapshot size (SELECT LIMIT), default 800 (unchanged, #2576 -- volume/order for ЕКБ stay byte-identical to pre-oblast behaviour). oblast_batch_size: int -- ДОПОЛНИТЕЛЬНАЯ reserved-квота для листингов области (#2576), default 100. Отдельный LIMIT, НЕ отъедает от batch_size ЕКБ -- гарантирует области честную обработку и одновременно не даёт всплеску свежих oblast-листингов вытеснить ЕКБ из top-N по scraped_at. budget_sec: float -- wall-clock budget per run, default 3600s. request_delay_sec: float -- delay between listings, default 6.0s. max_consecutive_blocks: int -- safety_min для BlockRatioBreaker (#3184): абортит прогоны короче окна, если ВСЕ попытки с начала были блоками (ни одного успеха), default 5. Основной критерий -- доля блоков в скользящем окне, см. settings.detail_backfill_block_ratio_window/ _threshold (module docstring). max_consecutive_failures: int -- порог обрыва по отказам-не-блокам, default 25 (см. комментарий у чтения параметра ниже). Lifecycle: update_heartbeat -> snapshot -> loop with budget guard -> mark_backfill_finished (done / banned при блоках / failed при нуле, #2674); mark_failed напрямую — только при исключении. """ batch_size = int(params.get("batch_size", 800)) oblast_batch_size = int(params.get("oblast_batch_size", 100)) budget_sec = float(params.get("budget_sec", 3600)) request_delay_sec = float(params.get("request_delay_sec", 6.0)) # #3184: теперь safety_min BlockRatioBreaker (см. module docstring) -- не общий # критерий обрыва, а страховка для прогонов короче окна. max_consecutive_blocks = int(params.get("max_consecutive_blocks", 5)) # Брейкер на отказы-НЕ-блоки. Блоки свой брейкер имели с самого начала, отказы — # нет, и это стоило трёх ночей подряд: 3-5 августа прогон делал ~1600 попыток, # получал 1600 отказов, ноль обогащений и выедал весь бюджет 9000 с (плюс 1600 # запросов через единственный прокси, #2638). Порог заметно выше блочного: пачка # мёртвых карточек (404 → ValueError в curl-режиме) не должна обрывать здоровый # прогон, а 25 отказов подряд без единого успеха — уже не невезение. max_consecutive_failures = int(params.get("max_consecutive_failures", 25)) 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'а. Без них один # зависший fetch_detail/rotate ронял весь run в zombie (run 423 завис 7.7ч). # Локальные float'ы (а не settings.* напрямую) — чтобы wait_for/format получали # число даже когда settings замокан MagicMock'ом в тестах. fetch_timeout_s = float(settings.avito_detail_fetch_timeout_s) rotate_timeout_s = float(settings.avito_proxy_rotate_settle_s) + 30.0 counters = AvitoDetailBackfillResult() current_counters: dict[str, int] = counters.to_dict() # Режим фетча: curl_cffi+backconnect (mproxy) когда use_curl=True; иначе браузер/auv. # avito_detail_backfill_use_curl=True переопределяет scraper_fetch_mode — detail_backfill # всегда идёт через backconnect чтобы не конкурировать за auv с SERP (city_sweep/full_load). use_curl = settings.avito_detail_backfill_use_curl browser_mode = settings.scraper_fetch_mode == "browser" and not use_curl session: AsyncSession | None = None browser_fetcher: BrowserFetcher | None = None own_session = False own_browser = False # kit AvitoScraper требует ScraperConfig позиционно (Strangler-инжекция #2133) — # RealScraperConfig проксирует settings.* так же, как читал legacy-конструктор без # аргументов (scraper_proxy_url и т.д. для _build_cffi_session()/_rotate_ip()). scraper = AvitoScraper(RealScraperConfig()) start = time.monotonic() logger.info( "avito_detail_backfill: run_id=%d mode=%s", run_id, "curl/backconnect" if use_curl else "browser", ) try: # Setup session (mirrors run_avito_pipeline lines 161-183). # use_curl=True: session=None + browser_fetcher=None → fetch_detail строит # эфемерную _build_detail_session() через settings.scraper_proxy_url (backconnect). # use_curl=False (legacy): browser_mode=True → BrowserFetcher/auv, как раньше. if browser_mode: # proxy_provider/use_pool/environment — обязательная часть проводки, # а не опция (#2698). Без них BrowserFetcher не кладёт "proxy" в тело # POST /fetch, сайдкар берёт свой env-прокси, и прогон уходит мимо # пула: ни выбора узла по affinity, ни учёта scrape_proxy_source_bans, # ни ротации при блоке. В логе это видно как proxy_lease_id=None. # # Ровно этот дефект уже чинили в house_imv_backfill (#2698) — там он # держал 35 из 35 отказов в каждом прогоне полтора месяца, пока # соседние свипы через ТОТ ЖЕ сайдкар тянули сотни объявлений. Здесь # он остался. # # Замер 27.08: прогон 5098 — mode=browser, proxy_lease_id=None, # 5 блоков подряд из 5 попыток, enriched=0. При этом пул здоров # (4 узла, все ok), а тот же URL через прокси отдаёт 200 и 3.3 МБ. # reuse_context=True (#3180/#3251): без него sidecar's browser.new_page() # создаёт НОВЫЙ изолированный context на КАЖДЫЙ /fetch — пройденный # QRATOR-PoW предыдущей карточки выбрасывается, и следующий запрос снова # холодный. У DomClick (#3118) это измеренно давало 100% блоков (26/26) # на изолированных контекстах против 5/5 успехов в тёплом контексте. # Сброс сожжённого context'а — ниже, через bf.request_context_reset() # (#3212: один раз за прогон, не на каждый блок — см. except-ветку). _cfg = RealScraperConfig() browser_fetcher = BrowserFetcher( source="avito", endpoint=settings.browser_http_endpoint, proxy_provider=RealProxyProvider(), use_pool=_cfg.use_proxy_pool_browser, # Без environment отказ «пул пуст» на этом пути мёртв — фетчер # молча ушёл бы на env-прокси сайдкара (#2616 шаг 1). environment=_cfg.environment, reuse_context=True, ) await browser_fetcher.__aenter__() own_browser = True 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 # config= обязателен: #2397 Part D1 — kit build_warmed_session() без него # молча теряет settings.scraper_proxy_url (sticky МГТС-прокси), см. модуль- # ный комментарий у импорта выше. session = await build_warmed_session(config=RealScraperConfig()) elif not use_curl: # curl_cffi legacy path (scraper_fetch_mode="curl_cffi", use_curl=False): # строим shared сессию через auv, как делает run_avito_city_sweep (kit). # Резолвер по источнику (#2825): пул scrape_proxies с учётом # scrape_proxy_source_bans, fallback на settings.scraper_proxy_url только # если пул пуст (легитимный dev/staging-сценарий). Пул не пуст, но все # забанены/нездоровы для avito -- resolve_proxy_url бросает # ProxyPoolExhaustedError (fail-closed, #2616): НАРОЧНО не ловим здесь -- # штатный except Exception ниже (mark_failed + logger.exception + raise) # уже даёт явную деградацию run'а с понятным логом, отдельный catch не нужен. own_session = True session = AsyncSession( impersonate=DEFAULT_IMPERSONATE, timeout=25, headers=DOCUMENT_HEADERS, proxies=http_proxies(resolve_proxy_url(db, "avito")), ) scraper._cffi = session runs_mod.update_heartbeat(db, run_id, current_counters) # SNAPSHOT: single SELECT at start -- NOT re-selected in loop. # Scope (#1814, расширено #2576): активные листинги ЕКБ + известных oblast- # городов (region 66). region_code на insert хардкодится в 66 (base.py) → # НЕ дискриминирует город; реальный признак города у Avito — путь URL # (/ekaterinburg/ для ЕКБ; legacy Москва/СПб/Тюмень — /moskva//sankt- # peterburg//tyumen/ — те по-прежнему вне scope, НЕ входят ни в ekb, ни в # oblast CTE). browser-fetch на legacy не-ЕКБ/не-oblast спотыкается → # curl-fallback → 429-бан curl-фингерпринта. Не тратим фетчи на мёртвые # (is_active) и на регионы вне scope. # # Два CTE вместо одного WHERE ... OR ...: ekb сохраняет ТОЧНО прежний # LIMIT/ORDER (#2576 требование "ЕКБ не деградирует") -- oblast НЕ может # вытеснить ЕКБ из batch_size ни при каком всплеске свежих oblast-строк # (ORDER BY ... scraped_at DESC в общем WHERE отдал бы приоритет самым # свежим независимо от города). oblast получает отдельную честную квоту # oblast_batch_size, добавленную ПОСЛЕ ekb-квоты (не вычтенную из неё). snapshot = ( db.execute( text( """ WITH ekb AS ( SELECT id, source_url, price_rub, 'ekb' AS city_scope FROM listings WHERE source = 'avito' AND detail_enriched_at IS NULL AND source_url IS NOT NULL AND is_active = TRUE AND source_url LIKE '%/ekaterinburg/%' -- сперва листинги без координат (#1967 — detail-страница -- даёт координаты здания), затем по свежести ORDER BY (lat IS NULL) DESC, scraped_at DESC NULLS LAST LIMIT CAST(:batch_size AS int) ), oblast AS ( SELECT id, source_url, price_rub, 'oblast' AS city_scope FROM listings WHERE source = 'avito' AND detail_enriched_at IS NULL AND source_url IS NOT NULL AND is_active = TRUE AND source_url LIKE ANY(CAST(:oblast_patterns AS text[])) ORDER BY (lat IS NULL) DESC, scraped_at DESC NULLS LAST LIMIT CAST(:oblast_batch_size AS int) ) SELECT id, source_url, price_rub, city_scope FROM ekb UNION ALL SELECT id, source_url, price_rub, city_scope FROM oblast """ ), { "batch_size": batch_size, "oblast_patterns": list(_OBLAST_AVITO_URL_PATTERNS), "oblast_batch_size": oblast_batch_size, }, ) .mappings() .all() ) if not snapshot: logger.info( "avito_detail_backfill: run_id=%d -- no pending listings " "(detail_enriched_at IS NULL = 0), done", run_id, ) runs_mod.mark_done(db, run_id, current_counters) return counters # #2576: разбивка ekb/oblast только для наблюдаемости -- .get() консервативен # (city_scope нет в mock-снапшотах старых тестов, дефолт "ekb" их не ломает). oblast_count = sum(1 for row in snapshot if row.get("city_scope") == "oblast") logger.info( "avito_detail_backfill: run_id=%d snapshot=%d (ekb=%d oblast=%d, " "budget=%.0fs delay=%.1fs max_blocks=%d mode=%s)", run_id, len(snapshot), len(snapshot) - oblast_count, oblast_count, budget_sec, request_delay_sec, max_consecutive_blocks, "curl/backconnect" if use_curl else "browser", ) # #3184: доля блоков в скользящем окне вместо голого "N подряд" -- см. module # docstring и app.services.backfill_block_breaker. breaker = BlockRatioBreaker( window_size=int(settings.detail_backfill_block_ratio_window), ratio_threshold=float(settings.detail_backfill_block_ratio_threshold), safety_min=max_consecutive_blocks, # snapshot_size гейтит safety-net (#3184 review MAJOR 2): пачка блоков в # начале ДЛИННОГО прогона не должна абортить его так же, как раньше -- # safety-net включён только когда снапшот короче окна и ratio-критерий # физически недостижим (см. backfill_block_breaker module docstring). snapshot_size=len(snapshot), ) consecutive_failures = 0 aborted_by_blocks = False abort_reason: str | None = None # #3288: прогон оборван пустым пулом прокси — «нечем ходить», а не бан. no_proxy_stop = False do_sleep = False items_since_warm = 0 # Счётчик попыток (успех ИЛИ отказ — оба тратят бюджет IP одинаково, см. # settings.avito_detail_backfill_rotate_after_attempts) на ТЕКУЩЕМ арендованном # прокси. Обнуляется на каждой ротации (успешной ИЛИ неуспешной — иначе # исчерпанный дневной лимит провайдера дёргал бы rotate_proxy на каждой попытке). attempts_since_rotation = 0 # #3251: сброс переиспользуемого browser-context'а (reuse_context=True выше) # разрешён РОВНО один раз за прогон — зеркалит domclick_detail_backfill (#3212). # Сброс на КАЖДЫЙ блок сам себя поддерживает: пройденный QRATOR-PoW живёт в # context'е, сброс его выбрасывает, повторная проверка с того же IP сразу # после принятой снова блокируется — одна осечка превращается в необратимый # каскад блоков (см. except-ветку ниже). context_reset_used = False # #3283g: сколько раз ЗА ПРОГОН уже сработала ротация-на-бан (в отличие от # context_reset_used эта ротация меняет IP вместе со сбросом, поэтому не # подвержена каскаду #3251 и может срабатывать несколько раз, но бюджетно). # Gap между любыми ротациями (по счётчику ИЛИ по бану) считается через # attempts_since_rotation — он и так обнуляется на каждой ротации. rotate_on_ban_used = 0 # Перепись причин (блоки + отказы) — переживает пересоздание контейнера, # в отличие от логов; см. _failure_signature. failure_census: Counter[str] = Counter() # #2764/#3178: диагнозы всех блоков прогона по ТИПУ исключения, С кратностями # (Counter, не set) — set терял их до решения: 5 прогонов подряд 4×platform+ # 1×infra и настоящий 2+2 приходили в mark_backfill_finished одинаково и # получали 'unknown' оба раза, хотя первый явно платформенный. Перепись # (kind -> count) решает _dominant_ban_kind в scrape_runs.py. block_ban_kinds: Counter[str] = Counter() for idx, row in enumerate(snapshot): # Budget guard elapsed = time.monotonic() - start if elapsed > budget_sec: logger.info( "avito_detail_backfill: run_id=%d -- budget %.0fs exhausted " "(elapsed=%.1fs), stopping at #%d/%d", run_id, budget_sec, elapsed, idx, len(snapshot), ) break # #1182 Phase 2: кооперативный SIGTERM-drain (деплой recreate scraper). # Чисто выходим на границе карточки — каждая карточка уже закоммичена, # mark_done с текущими счётчиками вызовется после loop'а. Snapshot — # pending-query (detail_enriched_at IS NULL), per-item идемпотентна → # следующий run сам до-резюмит остаток, resume_cursor не нужен. if shutdown_requested(): logger.info( "avito_detail_backfill: run_id=%d SIGTERM-drain — stopping at #%d/%d", run_id, idx, len(snapshot), ) break # Delay before each request except the first (джиттер ±30% — менее # роботизированный паттерн на органически прогретой сессии). if do_sleep: await asyncio.sleep(request_delay_sec * random.uniform(0.7, 1.4)) do_sleep = True source_url: str = row["source_url"] counters.attempted += 1 attempts_since_rotation += 1 # Normalise URL -> path (mirrors run_avito_city_sweep detail-phase, kit) 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: # config= обязателен здесь тоже (тот же footgun, что и в setup- # вызове build_warmed_session выше) -- иначе периодический re-warm # молча пересобрал бы сессию БЕЗ sticky-прокси после первого раза. new_session = await asyncio.wait_for( build_warmed_session(config=RealScraperConfig()), 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: # #1950: hard-timeout — зависший fetch (browser hang / curl-stall) не # должен блокировать loop навсегда (иначе budget-guard/heartbeat молчат # и run reaped как zombie). wait_for отменяет fetch → TimeoutError. # #3251: same-site SERP-якорь. Смысл он имеет ТОЛЬКО в browser- # режиме: сайдкар держит выдачу открытой якорной вкладкой и шлёт # её как Referer целевой навигации. На curl/own-session пути # сессия прогрета warm-batch'ем (см. referer= ниже, это ДРУГОЕ # поле и другой фетчер) — там origin остаётся None и fetch_detail # ведёт себя ровно как раньше. None также если source_url не # распарсился на город+категорию (см. docstring _serp_origin_for). serp_origin = _serp_origin_for(source_url) if browser_mode else None enrichment = await asyncio.wait_for( fetch_detail( item_url, cffi_session=session, browser_fetcher=browser_fetcher, referer=_AVITO_WARM_SEARCH_URL if use_curl else None, reconnect_on_block=not use_curl, # config= обязателен: kit fetch_detail без него читает # config.scraper_proxy_url=None для backconnect-gate'а # (elif not use_curl own-session path), теряя reconnect- # on-403 поведение legacy (settings.scraper_proxy_url # напрямую). Не влияет на use_curl=True (прод-дефолт) — # там reconnect_on_block=False уже гасит backconnect. config=RealScraperConfig(), origin=serp_origin, browser_referer=serp_origin, ), timeout=fetch_timeout_s, ) if save_detail_enrichment(db, enrichment): counters.enriched += 1 else: # #3338 (та же дыра, что #3332 у domclick): карточка взята и # разобрана, а UPDATE не задел ни одной строки — объявление # удалено/деактивировано между снимком и записью. Попытка была, # исхода не было: attempted переставал сходиться с # enriched + blocked + gone + failed, и расхождение читается как # потерянный отказ площадки. Исход failed: непрошедший UPDATE — не # успех, не блок и не gone (снятие метит is_active=FALSE сам, по 404). counters.failed += 1 # Печатаем ОБА идентификатора: WHERE в save ключуется по # source_id из разобранного HTML (item_id), а не по row-id из # снимка. При редиректе/подмене карточки строка listing_id # существует и жива — не нашлась строка с source_id=item_id. logger.warning( "avito_detail_backfill: run_id=%d listing %s -- карточка " "разобрана, но UPDATE не нашёл строку: listing_id=%s item_id=%s " "(WHERE по item_id из HTML)", run_id, source_url, row["id"], enrichment.item_id, ) if use_curl: items_since_warm += 1 breaker.record_success() consecutive_failures = 0 except AvitoListingGoneError as gone_exc: # #2034: мёртвый листинг (404 / removed) — НЕ блок, НЕ failed. # Координатные дыры в lat-null очереди в основном dead-листинги; # browser-mode рендерит их «Ошибка 404» без item-view → раньше это # ловилось как soft-block → consecutive-block breaker абортил run до # live-листингов (run 458: attempted=5 enriched=0 blocked=5 → abort). # 404 нейтрален к breaker'у (#3184: record_neutral -- ни окно, ни серия, # ни safety-net не трогаются). Метим is_active=FALSE → листинг уходит из # scope (snapshot SELECT фильтрует is_active = TRUE) и не тратит фетчи впредь. counters.gone += 1 breaker.record_neutral() # 404 — честный ответ площадки, значит тракт цел: серия отказов # прерывается (блочный брейкер 404 не трогает — см. #2034). consecutive_failures = 0 failure_census[_failure_signature(gone_exc)] += 1 try: with db.begin_nested(): db.execute( text("UPDATE listings SET is_active = FALSE WHERE id = :id"), {"id": row["id"]}, ) # #2674: 404 с площадки — самый достоверный сигнал снятия, # фиксируем его в дневной истории (listings_snapshots.status # был константой 'active' у всех строк, 394 299). Тот же # SAVEPOINT, что и UPDATE флага: снимок без флага (или # наоборот) невозможен. price_rub из snapshot-SELECT — # .get() консервативен ради mock-снапшотов старых тестов. gone_price = row.get("price_rub") if gone_price is not None: upsert_listing_snapshot( db, listing_id=row["id"], price_rub=gone_price, run_id=run_id, status="closed", ) except Exception: logger.warning( "avito_detail_backfill: run_id=%d failed to mark listing %s " "is_active=FALSE", run_id, row["id"], exc_info=True, ) logger.info( "avito_detail_backfill: run_id=%d GONE #%d/%d listing %s -> " "marked is_active=FALSE", run_id, counters.gone, len(snapshot), source_url, ) except (AvitoBlockedError, AvitoRateLimitedError) as e: # #3288: пустой пул — не блок и не отказ площадки: запрос не # уходил вовсе, следующая карточка упрётся ровно в то же самое. # Сюда он приезжает под видом блокировки, потому что fetch_detail # заворачивает в AvitoSidecarUnavailableError любое исключение # фетча. Исход честно failed (тот же разряд, что транспортные # сбои ниже) — тождество attempted = enriched + blocked + gone + # failed остаётся целым (#3338), а причину несёт no_proxy_stop=1 # и mark_failed с текстом про пул, как у домклика после #3283. if _caused_by_empty_pool(e): counters.failed += 1 logger.warning( "avito_detail_backfill: run_id=%d СТОП — пул прокси пуст, " "к площадке не ходили. enriched=%d attempted=%d", run_id, counters.enriched, counters.attempted, ) no_proxy_stop = True break ban_kind = ban_kind_of_exception(e) # #3288: в окно доли идёт только 'platform' — infra (отказ нашего # сайдкара) уходит в знаменатель, см. BlockRatioBreaker.record_block. breaker.record_block(ban_kind) counters.blocked += 1 failure_census[_failure_signature(e)] += 1 block_ban_kinds[ban_kind] += 1 do_sleep = False # #3251/#3283g: на настоящий бан ПЛОЩАДКОЙ (AvitoBlockedError и подтипы: # AvitoContentBlockedError, AvitoWarmupCookiesMissingError; НЕ на # AvitoSidecarUnavailableError — подтип AvitoRateLimitedError, отказ # НАШЕГО тракта, площадка ни при чём) есть два инструмента: # 1) rotate-on-ban (#3283g) — меняет exit-IP И сбрасывает context # (через _rotate_current_proxy), бюджетно (rotate_on_ban_max за # прогон) и с gap'ом (rotate_on_ban_min_gap попыток от последней # ротации любого происхождения, gap считаем по # attempts_since_rotation — она и так обнуляется на ротациях). # Смена IP лечит бан вживую за секунды (см. module docstring # задачи) и НЕ подвержена каскаду ниже — новый адрес не хранит # сожжённый PoW старого. # 2) bare context reset (#3251) — БЕЗ смены IP, ровно один раз за # прогон (domclick_detail_backfill #3212 — тот же приём). Сброс # на каждый блок сам себя поддерживает: пройденный QRATOR-PoW # живёт в context'е, сброс его выбрасывает → следующая проверка # с ТОГО ЖЕ IP снова блокируется — необратимый каскад. # Приоритет — rotate-on-ban: если она сработала, context уже сброшен # вместе со сменой IP, поэтому bare reset для этого события и всех # последующих (пока context_reset_used не сброшен — а он не сбрасывается # в рамках прогона) больше не нужен. Оба сброса в ОДНОМ событии не # допускаем: разные механизмы, но request_context_reset() внутри один, # и два подряд — то самое "два сброса подряд", которого просит избежать # задача, даже если один из них идёт с новым IP. if browser_fetcher is not None and isinstance(e, AvitoBlockedError): ban_budget = settings.avito_detail_backfill_rotate_on_ban_max if ban_budget > 0 and rotate_on_ban_used >= ban_budget: logger.info( "avito_detail_backfill: run_id=%d rotate-on-ban budget " "exhausted (%d/%d) -- staying on current proxy", run_id, rotate_on_ban_used, ban_budget, ) elif ( ban_budget > 0 and attempts_since_rotation < settings.avito_detail_backfill_rotate_on_ban_min_gap ): logger.info( "avito_detail_backfill: run_id=%d rotate-on-ban skipped -- " "only %d attempts since last rotation (need >=%d)", run_id, attempts_since_rotation, settings.avito_detail_backfill_rotate_on_ban_min_gap, ) elif ban_budget > 0: # Бюджет тратим и на неудачу — иначе отказавший rotate_proxy # дёргался бы на КАЖДОМ следующем бане до конца прогона. rotate_on_ban_used += 1 attempts_since_rotation = 0 logger.warning( "avito_detail_backfill: run_id=%d rotate-on-ban #%d/%d -- " "platform ban, rotating exit IP instead of bare reset", run_id, rotate_on_ban_used, ban_budget, ) # context_reset_used выставляем ТОЛЬКО по факту успеха: при # тихом отказе rotate_proxy (лимит, нет rotate_url, сеть) # сброса контекста внутри НЕ произошло, и съесть им # одноразовый bare-reset #3251 значило бы остаться и без # нового IP, и без сброса вообще. if await _rotate_current_proxy( db, run_id, browser_fetcher, trigger="rotate-on-ban" ): context_reset_used = True if not context_reset_used: context_reset_used = True browser_fetcher.request_context_reset() logger.warning( "avito_detail_backfill: run_id=%d BLOCKED #%d/%d (consecutive=%d): %s", run_id, idx + 1, len(snapshot), breaker.consecutive_blocks, e, ) # abort FIRST — на последнем блоке не тратим cooldown+research/rotate. # #3184: доля в скользящем окне ИЛИ safety-net (все попытки — блоки, ни # одного успеха, накопилось max_consecutive_blocks) -- см. module docstring. abort_reason = breaker.abort_reason() if abort_reason is not None: # Печатаем ту величину, по которой обрыв и произошёл. Раньше здесь # всегда стояло "%d consecutive blocks", хотя критерий с #3184 стал # ratio: прогон 5210 (14 блоков из 20, ровно порог) отпечатал # "ABORT -- 1 consecutive blocks" -- текущая серия в тот момент # действительно равнялась единице, но обрыв был не по ней. logger.warning( "avito_detail_backfill: run_id=%d ABORT -- %s, " "частая причина: %s. enriched=%d attempted=%d", run_id, breaker.abort_explanation(), _top_failure(failure_census) or "причина не определена", counters.enriched, counters.attempted, ) aborted_by_blocks = True 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 as e: # asyncio.wait_for → TimeoutError (py3.12: asyncio.TimeoutError — alias). # Ловим ДО общего Exception (TimeoutError ⊂ OSError ⊂ Exception). Зависший # fetch отменён → листинг failed, переходим к следующему (loop не зависает, # run не zombie #1950). Не считаем soft-блоком: rotate не дёргаем. counters.failed += 1 consecutive_failures += 1 breaker.record_failure() failure_census[_failure_signature(e)] += 1 logger.warning( "avito_detail_backfill: run_id=%d listing %s TIMEOUT (>%.0fs) -- skip", run_id, source_url, fetch_timeout_s, ) try: db.rollback() except Exception: pass except Exception as e: counters.failed += 1 consecutive_failures += 1 breaker.record_failure() failure_census[_failure_signature(e)] += 1 logger.warning( "avito_detail_backfill: run_id=%d listing %s failed: %s", run_id, source_url, e, ) try: db.rollback() except Exception: pass # Ротация exit-IP по счётчику попыток на арендованном прокси (задача поверх # #3298). browser_fetcher is not None -- см. _rotate_current_proxy docstring # за тем, почему только browser-режим держит lease с proxy_id. Обнуляем # счётчик ДО вызова (не после) -- неудача ротации не должна дёргать # rotate_proxy на КАЖДОЙ следующей попытке до конца прогона. if ( browser_fetcher is not None and attempts_since_rotation >= settings.avito_detail_backfill_rotate_after_attempts ): attempts_since_rotation = 0 await _rotate_current_proxy(db, run_id, browser_fetcher) if consecutive_failures >= max_consecutive_failures: logger.error( "avito_detail_backfill: run_id=%d ABORT -- %d отказов подряд без " "единого успеха, частая причина: %s. enriched=%d attempted=%d", run_id, consecutive_failures, _top_failure(failure_census) or "неизвестна", counters.enriched, counters.attempted, ) break if counters.attempted % 25 == 0: current_counters = counters.to_dict() runs_mod.update_heartbeat(db, run_id, current_counters) counters.duration_sec = time.monotonic() - start current_counters = counters.to_dict() # #3184: гистограмма длин пачек блоков (streak -> сколько раз встретилась) -- # иначе эффект правки на #2674-статистике нечем измерить постфактум. finalize() # досчитывает хвостовую пачку, если прогон оборвался посреди серии. streak_histogram = breaker.finalize() if streak_histogram: current_counters["block_streak_histogram"] = streak_histogram # type: ignore[assignment] # Какой из двух критериев оборвал прогон -- в counters, а не только в логах: # разбор простоя идёт SQL-запросом по scrape_runs, а не грепом контейнера, # и без этого ключа "banned" опять не отличить по причине (#3178). if abort_reason is not None: current_counters["abort_reason"] = abort_reason # type: ignore[assignment] if no_proxy_stop: # #3288 (как #3283 у домклика): остановка из-за пустого пула — НЕ блок, # поэтому и не aborted_by_blocks: иначе прогон уйдёт в 'banned' и запись # будет утверждать про площадку то, чего не было. Это отказ нашей стороны. current_counters["no_proxy_stop"] = 1 runs_mod.mark_failed( db, run_id, "пул прокси пуст — к площадке не ходили (#3288)", current_counters, ) logger.info( "avito_detail_backfill: run_id=%d FINISHED (пул пуст) -- attempted=%d " "enriched=%d blocked=%d gone=%d failed=%d duration=%.1fs", run_id, counters.attempted, counters.enriched, counters.blocked, counters.gone, counters.failed, counters.duration_sec, ) return counters runs_mod.mark_backfill_finished( db, run_id, current_counters, source="avito_detail_backfill", aborted_by_blocks=aborted_by_blocks, fail_hint=_top_failure(failure_census), ban_kinds=block_ban_kinds, ) logger.info( "avito_detail_backfill: run_id=%d FINISHED -- attempted=%d enriched=%d " "blocked=%d gone=%d failed=%d duration=%.1fs", run_id, counters.attempted, counters.enriched, counters.blocked, counters.gone, counters.failed, counters.duration_sec, ) return counters except Exception as exc: counters.duration_sec = time.monotonic() - start logger.exception( "avito_detail_backfill: run_id=%d FAILED after %.1fs", run_id, counters.duration_sec, ) current_counters = counters.to_dict() if _caused_by_empty_pool(exc): # #3384: пул был пуст ещё ДО первой карточки — lease берётся в # BrowserFetcher.__aenter__ (строка 427), поэтому NoProxyAvailableError # вылетает мимо стоп-механики цикла, которая и ставит no_proxy_stop. Без # ключа такой прогон (attempted=0, к площадке не ходили) неотличим от # любого другого падения: разбор простоя идёт SQL'ём по # counters.no_proxy_stop (#3288/#3367), а не грепом текста ошибки. current_counters["no_proxy_stop"] = 1 runs_mod.mark_failed(db, run_id, str(exc)[:1000], current_counters) raise finally: if own_session and session is not None: await session.close() if own_browser and browser_fetcher is not None: await browser_fetcher.__aexit__(None, None, None)