gendesign/tradein-mvp/backend/app/tasks/avito_detail_backfill.py
lekss361 01b5e73ea4
All checks were successful
Deploy Trade-In / changes (push) Successful in 12s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m33s
Deploy Trade-In / build-backend (push) Successful in 1m41s
Deploy Trade-In / deploy (push) Successful in 1m38s
feat(tradein): вся Свердловская область — 40 городов в city-sweep (#2879)
2026-08-13 19:08:44 +00:00

715 lines
40 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""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, так что передавать было бы нечего.
"""
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 DOCUMENT_HEADERS, http_proxies
from scraper_kit.providers.avito.detail import (
_AVITO_WARM_SEARCH_URL,
build_warmed_session,
fetch_detail,
research_in_session,
save_detail_enrichment,
)
from scraper_kit.providers.avito.serp import AvitoScraper
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.proxy_egress import resolve_proxy_url
from app.services.scraper_adapters import 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("<url>", 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())})"
@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 -- abort threshold, default 5.
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))
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:
browser_fetcher = BrowserFetcher(
source="avito", endpoint=settings.browser_http_endpoint
)
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="chrome120",
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",
)
consecutive_blocks = 0
consecutive_failures = 0
aborted_by_blocks = False
do_sleep = False
items_since_warm = 0
# Перепись причин (блоки + отказы) — переживает пересоздание контейнера,
# в отличие от логов; см. _failure_signature.
failure_census: Counter[str] = Counter()
# #2764: диагнозы всех блоков прогона по ТИПУ исключения. Сойдутся в один —
# он и попадёт в scrape_runs.ban_kind, разойдутся — 'unknown' (схлопывает
# mark_backfill_finished, один узел на все три backfill'а).
block_ban_kinds: set[str] = set()
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
# 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.
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(),
),
timeout=fetch_timeout_s,
)
if save_detail_enrichment(db, enrichment):
counters.enriched += 1
if use_curl:
items_since_warm += 1
consecutive_blocks = 0
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'у: consecutive_blocks НЕ трогаем (не растим
# и не сбрасываем). Метим is_active=FALSE → листинг уходит из scope
# (snapshot SELECT фильтрует is_active = TRUE) и не тратит фетчи впредь.
counters.gone += 1
# 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:
consecutive_blocks += 1
counters.blocked += 1
failure_census[_failure_signature(e)] += 1
block_ban_kinds.add(ban_kind_of_exception(e))
do_sleep = False
logger.warning(
"avito_detail_backfill: run_id=%d BLOCKED #%d/%d (consecutive=%d): %s",
run_id,
idx + 1,
len(snapshot),
consecutive_blocks,
e,
)
# abort FIRST — на последнем блоке не тратим cooldown+research/rotate
# (consecutive_blocks уже инкрементнут выше).
if consecutive_blocks >= max_consecutive_blocks:
logger.error(
"avito_detail_backfill: run_id=%d ABORT -- %d consecutive blocks, "
"IP rate-limited. enriched=%d attempted=%d",
run_id,
consecutive_blocks,
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
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
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
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()
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,
)
runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters.to_dict())
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)