"""pipeline.py — Avito sweep orchestrator (#2135 strangler, единственный в проде с #2397 Part E1). Изначально был strangler-COPY боевого `app.services.scrape_pipeline` (#2135); legacy-оригинал удалён (#2397 Part E1) — этот модуль теперь единственная orchestration-реализация. Развязка от `app.*` через инжекцию (см. `scraper_kit.contracts`): - `app.core.config.settings` → аргумент `config: ScraperConfig` - `app.services.matching` → аргумент `matcher: HouseMatcher` - `app.services.house_imv_backfill` → аргумент `enrichment: EnrichmentJobs` - `app.core.shutdown.shutdown_requested` → аргумент `shutdown_requested: Callable` - `app.services.scrape_runs` → локальный `scraper_kit.orchestration.runs` - провайдеры/база → `scraper_kit.providers.*` / `scraper_kit.base` Скопирована КРИТИЧНАЯ логика оркестрации (#2135): - `run_avito_pipeline` — одиночный anchor (SEARCH→SAVE→GROUP→HOUSES→DETAIL) - `run_avito_city_sweep` — city-sweep + ban/rotation + partial-intake + IMV (F1) - `run_avito_newbuilding_sweep`— citywide novostroyka-обход + save (F2) - `run_yandex_city_sweep` — combos (room×price) + save + address-enrich (F2) - `run_cian_city_sweep` — anchor SERP → detail(+price-hist) → newbuilding (F2) - `run_cian_full_load` — region-wide bisection + resume/checkpoint (F2) - `run_yandex_full_load` — region-wide bisection + price-history (F2) - `run_avito_full_load` — citywide bisection (exhaustive/incremental) (F2) - `run_domclick_city_sweep` — BFF citywide sweep + честный статус (F2) Ban-rotation state-machine сведена в единый хелпер `_try_rotate_within_budget` (в оригинале инлайн-дублировалась в house/detail/house-detail/enrich фазах). Yandex sweep'ы дёргают продуктовую логику (price-history + address-parse) через инжектируемый `EnrichmentJobs` (record_yandex_price_history / extract_address_from_title / address_has_house_number) вместо прямых импортов app.services.yandex_*. """ from __future__ import annotations import asyncio import logging import random from collections.abc import Callable from contextlib import AsyncExitStack from dataclasses import dataclass, field, fields from datetime import date, timedelta from typing import TYPE_CHECKING, Any from urllib.parse import urlparse from curl_cffi.requests import AsyncSession from sqlalchemy import text from scraper_kit.avito_exceptions import AvitoBlockedError, AvitoRateLimitedError from scraper_kit.base import ScrapedLot, save_listings from scraper_kit.browser_fetcher import BrowserFetcher from scraper_kit.orchestration import runs from scraper_kit.providers.avito.detail import fetch_detail, save_detail_enrichment from scraper_kit.providers.avito.houses import ( fetch_house_catalog, save_house_catalog_enrichment, ) from scraper_kit.providers.avito.serp import AvitoScraper from scraper_kit.providers.cian.detail import fetch_detail as cian_fetch_detail from scraper_kit.providers.cian.detail import save_detail_enrichment as cian_save_detail_enrichment from scraper_kit.providers.cian.newbuilding import ( fetch_newbuilding, save_newbuilding_enrichment, ) from scraper_kit.providers.cian.serp import CianScraper from scraper_kit.providers.domclick.serp import ROOM_BUCKETS, DomClickScraper from scraper_kit.providers.yandex.serp import ( DEFAULT_PRICE_RANGES, ROOM_PATH, YandexRealtyScraper, ) if TYPE_CHECKING: from sqlalchemy.orm import Session from scraper_kit.contracts import ( EnrichmentJobs, HouseMatcher, ProxyProvider, ScraperConfig, ) logger = logging.getLogger(__name__) # Регион по умолчанию для listings.region_code (в старом save_listings хардкод 66). # Вынесен параметром run_*-функций — kit не знает про конкретный регион. DEFAULT_REGION_CODE: int = 66 def _max_rotations(config: ScraperConfig, source: str) -> int: """Лимит IP-ротаций по провайдеру (avito/cian/yandex).""" if source == "cian": return config.cian_proxy_max_rotations if source == "yandex": return config.yandex_proxy_max_rotations return config.avito_proxy_max_rotations async def _rotate_proxy_ip( config: ScraperConfig, *, reason: str, rotations_done: int, source: str = "avito", ) -> bool: """Вызвать changeip-ссылку mobileproxy и подождать settle (#1790/#1848). Provider-aware: выбирает rotate_url и max_rotations по параметру `source` (avito / cian / yandex). Fallback-цепочка для rotate_url: avito: avito_proxy_rotate_url cian: cian_proxy_rotate_url → avito_proxy_rotate_url yandex: yandex_proxy_rotate_url → avito_proxy_rotate_url Аналог AvitoScraper._rotate_ip, но на уровне pipeline — используется в enrichment-фазах где нет доступа к экземпляру scraper'а. Returns True при успешной смене IP, False если rotate_url не задан или ошибка. Логирует каждую ротацию (причина, порядковый номер, source, остаток лимита). """ if source == "cian": rotate_url = config.cian_proxy_rotate_url or config.avito_proxy_rotate_url max_rot = config.cian_proxy_max_rotations elif source == "yandex": rotate_url = config.yandex_proxy_rotate_url or config.avito_proxy_rotate_url max_rot = config.yandex_proxy_max_rotations else: # avito (default) — backward compat rotate_url = config.avito_proxy_rotate_url max_rot = config.avito_proxy_max_rotations if not rotate_url: return False sep = "&" if "?" in rotate_url else "?" # #1950: retry changeip-GET короткими попытками вместо одношотного 30s-timeout. # Если changeip-сервер завис на 30s → весь прогон ждёт зря + ротация считается # успешной (False возвращался). Теперь: proxy_rotate_attempts попыток по # proxy_rotate_attempt_timeout_s каждая; возвращаем True при первом успехе. _attempts = config.proxy_rotate_attempts _attempt_timeout = config.proxy_rotate_attempt_timeout_s for attempt in range(_attempts): try: async with AsyncSession(timeout=_attempt_timeout) as rot: resp = await rot.get(f"{rotate_url}{sep}format=json") new_ip: str = "" try: data = resp.json() new_ip = str(data.get("new_ip", "")) except Exception: pass await asyncio.sleep(config.avito_proxy_rotate_settle_s) logger.info( "pipeline: IP rotated via changeip — source=%s reason=%s " "rotation=#%d attempt=%d/%d remaining=%d new_ip=%s", source, reason, rotations_done + 1, attempt + 1, _attempts, max_rot - rotations_done - 1, new_ip or "unknown", ) return True except Exception: logger.warning( "pipeline: IP rotation attempt %d/%d failed " "(source=%s reason=%s rotation=#%d)", attempt + 1, _attempts, source, reason, rotations_done + 1, exc_info=True, ) if attempt < _attempts - 1: await asyncio.sleep(1.0) logger.error( "pipeline: all %d changeip attempts failed " "(source=%s reason=%s rotation=#%d)", _attempts, source, reason, rotations_done + 1, ) return False async def _try_rotate_within_budget( config: ScraperConfig, *, reason: str, rotations_done: int, source: str = "avito", ) -> tuple[bool, int]: """Единый шаг ban-rotation state-machine (свод инлайн-дублей #2135). Вызывается при достижении порога ПОСЛЕДОВАТЕЛЬНЫХ блоков. Если бюджет ротаций ещё не исчерпан (`rotations_done < max_rotations(source)`) — дёргает changeip и увеличивает счётчик ротаций. Возвращает `(rotated, rotations_done')`: - rotated=True → IP сменён; вызывающий сбрасывает consecutive и продолжает. - rotated=False → бюджет исчерпан ЛИБО changeip не удался; вызывающий делает abort/propagate (листинги уже сохранены). Поведение идентично прежним инлайн-веткам: инкремент rotations_done происходит только при наличии бюджета (когда реально пытаемся ротировать). """ if rotations_done < _max_rotations(config, source): rotated = await _rotate_proxy_ip( config, reason=reason, rotations_done=rotations_done, source=source ) return rotated, rotations_done + 1 return False, rotations_done # Per-anchor watchdog timeout (секунды). Если обработка одного anchor'а занимает # дольше этого значения (зависший HTTP-запрос, прокси hang), asyncio.wait_for # прерывает её с TimeoutError, sweep продолжается со следующим anchor'ом. # Диапазон 180-300 с; 240 — разумный середина (4 мин > любой нормальный anchor). ANCHOR_TIMEOUT_SEC: int = 240 # Константы для расчёта watchdog-таймаута Avito city sweep. # Avito detail-fetch идёт через браузерный сервис с page.goto timeout 60s, и # anti-bot страницы нередко доходят до этого лимита. Фиксированный ANCHOR_TIMEOUT_SEC=240 # guillotine'ит anchor ещё в detail-фазе → SERP-счётчики теряются, run показывает # lots_fetched=0 при реально сохранённых листингах. Таймаут масштабируется: # _avito_anchor_timeout = ANCHOR_TIMEOUT_SEC # + detail_top_n × _AVITO_PER_DETAIL_S # + (_AVITO_HOUSES_BUDGET_S if enrich_houses else 0) # Пример: detail_top_n=20 + houses → 240 + 1000 + 180 = 1420s/anchor (worst-case). _AVITO_PER_DETAIL_S: float = 50.0 # секунд на один detail-fetch (incl. occasional 60s timeout) _AVITO_HOUSES_BUDGET_S: float = 180.0 # бюджет на house-enrich фазу # Fail-fast: если detail phase накапливает подряд N timeout/ошибок, прерываем её # (аналог consecutive_blocks для AvitoBlockedError в step 5). _AVITO_DETAIL_CONSECUTIVE_TIMEOUT_ABORT: int = 5 # #1820: порог ПОСЛЕДОВАТЕЛЬНЫХ блоков (AvitoBlockedError/AvitoRateLimitedError) внутри # enrichment-фаз run_avito_pipeline. Один заблокированный фетч (house/detail) больше НЕ # роняет весь прогон — блок ловится пер-item (как обычная ошибка), и лишь N подряд # прерывают enrichment-фазу. Listings из SEARCH+SAVE сохраняются всегда. _AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT: int = 3 # #2160: константы для расчёта watchdog-таймаута Cian city sweep. При # USE_PROXY_POOL_BROWSER=true каждый SERP-фетч идёт через camoufox с relaunch при смене # прокси (page.goto timeout 60s + overhead) = 13-45s/страница, а якорь cian = 4 room-buckets # × pages_per_anchor SERP-фетчей + detail_top_n detail-фетчей + houses-enrich. Фиксированный # ANCHOR_TIMEOUT_SEC=240 гильотинит якорь ещё в SERP-фазе → save_listings не успевает, # lots_fetched=0, после N подряд → mark_banned. Таймаут масштабируется: # _cian_anchor_timeout = max( # ANCHOR_TIMEOUT_SEC, # _CIAN_ROOM_BUCKETS × pages_per_anchor × (request_delay_sec + _CIAN_PER_SERP_FETCH_S) # + detail_top_n × _CIAN_PER_DETAIL_S # + (_CIAN_HOUSES_BUDGET_S if enrich_houses else 0)) # Пример (дефолты): 4×3×(5+70) + 10×50 + 180 = 900 + 500 + 180 = 1580s/anchor (worst-case, # это watchdog от зависания, не ожидаемая длительность). _CIAN_ROOM_BUCKETS: int = 4 # fetch_around_multi_room итерирует ((1,),(2,),(3,),(4,)), # см. providers/cian/serp.py:215 _CIAN_PER_SERP_FETCH_S: float = 70.0 # browser fetch через пул: page.goto 60s + camoufox relaunch _CIAN_PER_DETAIL_S: float = 50.0 # секунд на один detail-fetch (как avito) _CIAN_HOUSES_BUDGET_S: float = 180.0 # бюджет на house-enrich фазу (как avito) def _cian_anchor_timeout_s( pages_per_anchor: int, request_delay_sec: float, detail_top_n: int, enrich_houses: bool, ) -> float: """Масштабируемый watchdog-таймаут одного Cian anchor'а (секунды). Чистая функция для тестируемости формулы (#2160). Никогда не опускается ниже ANCHOR_TIMEOUT_SEC (max-ветка): при малых параметрах фиксированный минимум остаётся. """ return max( float(ANCHOR_TIMEOUT_SEC), _CIAN_ROOM_BUCKETS * pages_per_anchor * (request_delay_sec + _CIAN_PER_SERP_FETCH_S) + detail_top_n * _CIAN_PER_DETAIL_S + (_CIAN_HOUSES_BUDGET_S if enrich_houses else 0.0), ) # Default anchors ЕКБ — 5 точек покрытия города EKB_ANCHORS: list[tuple[float, float, str]] = [ (56.8400, 60.6050, "Центр"), (56.7950, 60.5300, "ЮЗ"), (56.8970, 60.6100, "Уралмаш"), (56.7700, 60.5500, "Академический"), (56.8650, 60.6200, "Пионерский"), ] # Anchors для oblast-городов вне ЕКБ (B1 rollout — Свердловская область, region 66). # По 1 anchor на город (центр) — radius_m (передаётся отдельно, в scrape_schedules. # default_params) должен покрывать город целиком; в отличие от EKB_ANCHORS (5 точек # на большой город), эти города компактнее — одного anchor'а достаточно. # Domclick (BFF, city_id-based, НЕ anchors-based) сюда не входит — отдельный B2 rollout. CITY_ANCHORS: dict[str, list[tuple[float, float, str]]] = { "nizhniy_tagil": [(57.910, 59.980, "Н.Тагил центр")], "kamensk_uralskiy": [(56.414, 61.918, "Каменск центр")], "pervouralsk": [(56.908, 59.943, "Первоуральск центр")], "verkhnyaya_pyshma": [(56.976, 60.578, "В.Пышма центр")], "serov": [(59.604, 60.578, "Серов центр")], } def get_city_anchors(city_slug: str | None) -> list[tuple[float, float, str]] | None: """Anchors для oblast-города по slug (CITY_ANCHORS). Неизвестный slug/None → None. None — сигнал caller'у (scheduler.py _job_*_city_sweep) падать обратно на EKB_ANCHORS (run_avito_city_sweep/run_cian_city_sweep) либо на единственный центральный anchor ЕКБ (run_yandex_city_sweep combos-mode) — см. `anchors if anchors is not None else EKB_ANCHORS` в соответствующих run_*_city_sweep. """ if city_slug is None: return None return CITY_ANCHORS.get(city_slug) _CHROME_HEADERS = { "Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8", "Accept-Language": "ru-RU,ru;q=0.9,en;q=0.8", "Cache-Control": "max-age=0", "Sec-Fetch-Dest": "document", "Sec-Fetch-Mode": "navigate", "Sec-Fetch-Site": "none", "Sec-Fetch-User": "?1", "Upgrade-Insecure-Requests": "1", } def _avito_proxies(config: ScraperConfig) -> dict[str, str] | None: """#623/#806: mobile-proxy egress для shared-сессии sweep'а. КРИТИЧНО: sweep создаёт собственную shared AsyncSession и присваивает её `scraper._cffi`, перезатирая проксированную сессию из `AvitoScraper.__aenter__`. Через эту же сессию идёт и SERP, и detail-энричмент (fetch_detail). Без прокси detail-страницы летят с datacenter-IP → HTTP 429 → весь run mark_banned. Пусто (env не задан) → прямое подключение (dev/no-op). Читает config.scraper_proxy_url (SCRAPER_PROXY_URL > AVITO_PROXY_URL fallback). """ url = config.scraper_proxy_url return {"http": url, "https": url} if url else None # ── Result dataclasses ────────────────────────────────────────── @dataclass class PipelineCounters: """Counters одного pipeline run для логов/админки.""" lots_fetched: int = 0 lots_inserted: int = 0 lots_updated: int = 0 unique_houses: int = 0 houses_enriched: int = 0 houses_failed: int = 0 detail_attempted: int = 0 detail_enriched: int = 0 detail_failed: int = 0 errors: list[str] = field(default_factory=list) @dataclass class PipelineResult: """Полный результат full Avito pipeline.""" anchor_lat: float anchor_lon: float radius_m: int counters: PipelineCounters enrich_houses: bool enrich_detail_top_n: int touched_house_ids: set[int] = field(default_factory=set) # ── Main orchestrator ─────────────────────────────────────────── async def run_avito_pipeline( db: Session, lat: float, lon: float, radius_m: int = 1500, *, config: ScraperConfig, matcher: HouseMatcher, enrich_houses: bool = True, enrich_detail_top_n: int = 10, pages: int = 1, request_delay_sec: float | None = None, shared_session: AsyncSession | None = None, shared_browser: BrowserFetcher | None = None, proxy_provider: ProxyProvider | None = None, region_code: int = DEFAULT_REGION_CODE, ) -> PipelineResult: """Full Avito search → houses → detail enrichment pipeline. Steps: 1. SEARCH: AvitoScraper().fetch_around(lat, lon, radius_m, pages=pages) 2. SAVE: save_listings(db, lots) → (inserted, updated) 3. GROUP: dedupe house_url из новых lots → unique house path set 4. ENRICH_HOUSES: fetch_house_catalog + save per house (with sleep between) 5. ENRICH_DETAIL: top-N listings → fetch_detail + save (abort on 3 blocks) shared_session — если передан, используется без закрытия (lifecycle у вызывающего). shared_browser — если передан в browser-mode, используется без закрытия. Блокировки (AvitoBlockedError / AvitoRateLimitedError) в SEARCH-фазе (шаг 1) по-прежнему propagate наружу (без листингов прогон бессмыслен). В enrichment-фазах (houses шаг 4 + detail шаг 5 + house-detail шаг 5b) блок ловится ПЕР-ITEM как обычная ошибка (#1820): один заблокированный фетч НЕ роняет весь прогон. Лишь _AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT блоков ПОДРЯД прерывают enrichment — но PipelineResult всё равно возвращается с уже собранными listings/обогащениями. """ counters = PipelineCounters() lots: list[ScrapedLot] = [] detail_delay = request_delay_sec if request_delay_sec is not None else 7.0 # #1820: общий счётчик ПОСЛЕДОВАТЕЛЬНЫХ блоков на все enrichment-фазы (4/5/5b). # Сбрасывается на любом успешном fetch; при достижении порога — ротация IP (#1790) # или abort если бюджет ротаций исчерпан. enrich_consecutive_blocks = 0 enrich_rotations_done = 0 # #1790: сколько раз уже сменили IP в этом run enrichment_aborted = False browser_mode = config.scraper_fetch_mode == "browser" session: AsyncSession | None = None browser_fetcher: BrowserFetcher | None = None own_session = False own_browser = False scraper = AvitoScraper(config) if browser_mode: browser_fetcher = shared_browser if browser_fetcher is None: browser_fetcher = BrowserFetcher( source="avito", endpoint=config.browser_http_endpoint, proxy_provider=proxy_provider, use_pool=config.use_proxy_pool_browser, ) await browser_fetcher.__aenter__() own_browser = True scraper._browser = browser_fetcher else: own_session = shared_session is None session = shared_session or AsyncSession( impersonate="chrome120", timeout=25, headers=_CHROME_HEADERS, proxies=_avito_proxies(config), ) scraper._cffi = session try: # ── Step 1: search ────────────────────────────────────── try: lots = await scraper.fetch_around( lat, lon, radius_m, pages=pages, delay_override_sec=request_delay_sec, ) counters.lots_fetched = len(lots) logger.info( "pipeline:search anchor=(%.4f,%.4f) r=%dm fetched=%d", lat, lon, radius_m, len(lots), ) except (AvitoBlockedError, AvitoRateLimitedError): logger.error("pipeline:search BLOCKED by Avito at anchor=(%.4f,%.4f)", lat, lon) raise except Exception as e: logger.exception("pipeline:search failed") counters.errors.append(f"search: {e}") return PipelineResult(lat, lon, radius_m, counters, enrich_houses, enrich_detail_top_n) # ── Step 2: save listings ─────────────────────────────── if lots: try: counters.lots_inserted, counters.lots_updated = save_listings( db, lots, matcher=matcher, region_code=region_code ) except Exception as e: logger.exception("pipeline:save_listings failed") counters.errors.append(f"save_listings: {e}") # ── Step 3: group by house ────────────────────────────── unique_house_paths: set[str] = set() for lot in lots: if lot.house_url: try: parsed = urlparse(lot.house_url) path = parsed.path if parsed.path else lot.house_url unique_house_paths.add(path) except Exception: continue counters.unique_houses = len(unique_house_paths) logger.info("pipeline:group_by_house unique=%d", counters.unique_houses) # ── Step 4: enrich houses с shared session + sleep ────── touched_house_ids: set[int] = set() if enrich_houses and unique_house_paths and not enrichment_aborted: house_paths_list = list(unique_house_paths) for idx, house_path in enumerate(house_paths_list): try: enrichment = await fetch_house_catalog( house_path, cffi_session=session, browser_fetcher=browser_fetcher ) house_counters = save_house_catalog_enrichment(db, enrichment, matcher=matcher) counters.houses_enriched += 1 enrich_consecutive_blocks = 0 hid = house_counters.get("house_id") if hid: touched_house_ids.add(int(hid)) except (AvitoBlockedError, AvitoRateLimitedError) as e: # #1820: блок ловим пер-item как обычную ошибку — НЕ роняем прогон. enrich_consecutive_blocks += 1 counters.houses_failed += 1 counters.errors.append(f"BLOCKED_HOUSE {house_path}: {e}") logger.warning( "pipeline:house BLOCKED at %s (consecutive=%d)", house_path, enrich_consecutive_blocks, ) if enrich_consecutive_blocks >= _AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT: # #1790: Перед abort'ом пробуем сменить IP (changeip). # После ротации сбрасываем счётчик и продолжаем — НЕ abort. # Если ротации исчерпаны — abort enrichment (listings уже сохранены). rotated, enrich_rotations_done = await _try_rotate_within_budget( config, reason=f"house consecutive={enrich_consecutive_blocks}", rotations_done=enrich_rotations_done, ) if rotated: enrich_consecutive_blocks = 0 logger.info( "pipeline:enrich ROTATE — IP changed, " "reset consecutive blocks (rotation #%d/%d)", enrich_rotations_done, config.avito_proxy_max_rotations, ) continue logger.error( "pipeline:enrich ABORT — %d consecutive blocks (IP rate-limited), " "listings preserved", enrich_consecutive_blocks, ) counters.errors.append( f"AVITO_BLOCKED — enrichment aborted after " f"{enrich_consecutive_blocks} consecutive blocks (house phase)" ) enrichment_aborted = True break except Exception as e: logger.warning("pipeline:house_enrich failed for %s: %s", house_path, e) counters.houses_failed += 1 counters.errors.append(f"house {house_path}: {e}") # #1368: save_house_catalog_enrichment не делает внутренний rollback, # поэтому упавший промежуточный statement оставляет общий Session в # InFailedSqlTransaction. Без сброса — каскад «current transaction is # aborted» на все оставшиеся дома/listings этого якоря. try: db.rollback() except Exception: pass if idx < len(house_paths_list) - 1: await asyncio.sleep(detail_delay) # ── Step 5: enrich detail для top-N priority listings ─── if enrich_detail_top_n > 0 and not enrichment_aborted: priority_rows = ( db.execute( text(""" SELECT source_url FROM listings WHERE source = 'avito' AND source_url IS NOT NULL AND ( ( detail_enriched_at IS NULL AND price_rub > 0 AND ST_DWithin( geom::geography, ST_MakePoint(:lon, :lat)::geography, :radius ) ) OR ( detail_enriched_at IS NULL AND scraped_at > NOW() - INTERVAL '2 hours' ) ) ORDER BY scraped_at DESC NULLS LAST LIMIT :limit """), { "lat": lat, "lon": lon, "radius": radius_m * 2, "limit": enrich_detail_top_n, }, ) .mappings() .all() ) detail_house_paths: set[str] = set() for idx, row in enumerate(priority_rows): source_url: str = row["source_url"] counters.detail_attempted += 1 try: item_url = ( urlparse(source_url).path if source_url.startswith("http") else source_url ) enrichment_detail = await fetch_detail( item_url, cffi_session=session, browser_fetcher=browser_fetcher ) if save_detail_enrichment(db, enrichment_detail): counters.detail_enriched += 1 # #871: house-link есть в detail SSR (в SERP-карточке — JS-only). # Копим catalog-пути для house-enrich ниже (Step 5b). if enrichment_detail.house_catalog_url: hp = urlparse(enrichment_detail.house_catalog_url).path if hp: detail_house_paths.add(hp) enrich_consecutive_blocks = 0 except (AvitoBlockedError, AvitoRateLimitedError) as e: # #1820: блок ловим пер-item, НЕ роняем прогон. Счётчик общий с house-фазой. enrich_consecutive_blocks += 1 counters.detail_failed += 1 counters.errors.append(f"BLOCKED_DETAIL #{idx}: {e}") logger.warning( "pipeline:detail BLOCKED #%d/%d (consecutive=%d): %s", idx + 1, len(priority_rows), enrich_consecutive_blocks, e, ) if enrich_consecutive_blocks >= _AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT: # #1790: Перед abort'ом пробуем сменить IP (changeip). # После ротации сбрасываем счётчик — продолжаем detail-обход. # Если ротации исчерпаны — abort detail-фазы (listings уже сохранены). rotated, enrich_rotations_done = await _try_rotate_within_budget( config, reason=f"detail consecutive={enrich_consecutive_blocks}", rotations_done=enrich_rotations_done, ) if rotated: enrich_consecutive_blocks = 0 logger.info( "pipeline:enrich ROTATE — IP changed, " "reset consecutive blocks (rotation #%d/%d)", enrich_rotations_done, config.avito_proxy_max_rotations, ) continue logger.error( "pipeline:detail ABORT — %d consecutive blocks (IP rate-limited), " "listings preserved", enrich_consecutive_blocks, ) counters.errors.append( f"AVITO_BLOCKED — enrichment aborted after " f"{idx + 1}/{len(priority_rows)} details (consecutive blocks)" ) # #1820: НЕ raise — прерываем только detail-фазу (и Step 5b), # PipelineResult возвращается с уже сохранёнными listings/houses. enrichment_aborted = True break except Exception as e: counters.detail_failed += 1 counters.errors.append(f"detail {source_url}: {e}") logger.warning("pipeline:detail_enrich failed for %s: %s", source_url, e) # #1386: НЕ сбрасываем consecutive_blocks здесь. Транзиентная # сетевая/parse-ошибка (5xx, timeout, connection reset) — правдоподобный # побочный эффект активного IP-блока (Avito не отдаёт 403/429 единообразно). # Сброс маскировал бы продолжающийся блок (block, block, timeout→reset, ...) # и не давал бы сработать early-abort на 3 consecutive blocks. # Сброс на успехе (см. выше) остаётся. # #1368: save_detail_enrichment без внутреннего rollback мог оставить # общий Session в aborted-state — сбрасываем, иначе каскад на оставшиеся # detail-строки и Step 5b. try: db.rollback() except Exception: pass if idx < len(priority_rows) - 1: jitter = random.uniform(0.8, 1.2) await asyncio.sleep(detail_delay * jitter) # ── Step 5b: enrich houses, найденные через detail-страницы (#871) ── # SERP-карточки больше не несут house-link (JS-render) → дома берём из # detail SSR. Дедуп против обработанных в Step 4 (unique_house_paths). new_house_paths = detail_house_paths - unique_house_paths if enrich_houses and new_house_paths and not enrichment_aborted: counters.unique_houses += len(new_house_paths) nh_list = list(new_house_paths) for h_idx, house_path in enumerate(nh_list): try: enrichment = await fetch_house_catalog( house_path, cffi_session=session, browser_fetcher=browser_fetcher ) house_counters = save_house_catalog_enrichment( db, enrichment, matcher=matcher ) counters.houses_enriched += 1 enrich_consecutive_blocks = 0 hid = house_counters.get("house_id") if hid: touched_house_ids.add(int(hid)) except (AvitoBlockedError, AvitoRateLimitedError) as e: # #1820: блок пер-item, общий счётчик; не роняем прогон. enrich_consecutive_blocks += 1 counters.houses_failed += 1 counters.errors.append(f"BLOCKED_HOUSE(detail) {house_path}: {e}") logger.warning( "pipeline:house(detail) BLOCKED at %s (consecutive=%d)", house_path, enrich_consecutive_blocks, ) if enrich_consecutive_blocks >= _AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT: # #1790: попытка ротации перед abort'ом (общий бюджет run). rotated, enrich_rotations_done = await _try_rotate_within_budget( config, reason=f"house-detail consecutive={enrich_consecutive_blocks}", rotations_done=enrich_rotations_done, ) if rotated: enrich_consecutive_blocks = 0 logger.info( "pipeline:enrich ROTATE — IP changed, " "reset consecutive blocks (rotation #%d/%d)", enrich_rotations_done, config.avito_proxy_max_rotations, ) continue logger.error( "pipeline:enrich ABORT — %d consecutive blocks " "(IP rate-limited), listings preserved", enrich_consecutive_blocks, ) counters.errors.append( f"AVITO_BLOCKED — enrichment aborted after " f"{enrich_consecutive_blocks} consecutive blocks " f"(house-detail phase)" ) enrichment_aborted = True break except Exception as e: logger.warning( "pipeline:house(detail) enrich failed for %s: %s", house_path, e ) counters.houses_failed += 1 counters.errors.append(f"house(detail) {house_path}: {e}") # #1368: см. Step 4 — сбрасываем aborted-state общего Session, # иначе каскад на оставшиеся дома detail-фазы. try: db.rollback() except Exception: pass if h_idx < len(nh_list) - 1: await asyncio.sleep(detail_delay) logger.info( "pipeline:done anchor=(%.4f,%.4f) lots=%d (ins=%d/upd=%d) " "houses=%d/%d detail=%d/%d touched_house_ids=%d errors=%d", lat, lon, counters.lots_fetched, counters.lots_inserted, counters.lots_updated, counters.houses_enriched, counters.unique_houses, counters.detail_enriched, counters.detail_attempted, len(touched_house_ids), len(counters.errors), ) return PipelineResult( lat, lon, radius_m, counters, enrich_houses, enrich_detail_top_n, touched_house_ids=touched_house_ids, ) 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) @dataclass class CitySweepCounters: """Aggregate counters для full city sweep run.""" anchors_total: int = 0 anchors_done: int = 0 lots_fetched: int = 0 lots_inserted: int = 0 lots_updated: int = 0 unique_houses: int = 0 houses_enriched: int = 0 houses_failed: int = 0 detail_attempted: int = 0 detail_enriched: int = 0 detail_failed: int = 0 imv_attempted: int = 0 imv_enriched: int = 0 imv_failed: int = 0 errors_count: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} async def run_avito_city_sweep( db: Session, *, run_id: int, config: ScraperConfig, matcher: HouseMatcher, enrichment: EnrichmentJobs, shutdown_requested: Callable[[], bool] = lambda: False, proxy_provider: ProxyProvider | None = None, anchors: list[tuple[float, float, str]] | None = None, radius_m: int = 1500, pages_per_anchor: int = 3, enrich_houses: bool = True, detail_top_n: int = 10, request_delay_sec: float = 7.0, enrich_imv: bool = True, region_code: int = DEFAULT_REGION_CODE, ) -> CitySweepCounters: """Full city sweep: iterate anchors × pages → save → enrich houses + detail → IMV. - Single shared AsyncSession на весь sweep (один TLS fingerprint) - AvitoBlockedError/AvitoRateLimitedError → ABORT sweep + mark_banned (status='banned') - Прочие errors per-anchor логируются, не валят весь sweep - После всех anchor'ов: если enrich_imv=True — IMV-оценка тронутых домов (cooperative cancel + per-house graceful error handling) SERP-счётчики (lots_fetched/inserted/updated) записываются в sweep-level counters НЕМЕДЛЕННО после save_listings — до начала detail/houses фазы. Даже если anchor превышает watchdog-таймаут в detail-фазе, SERP-результаты не теряются. Инжекция (#2135): config/matcher/enrichment/shutdown_requested приходят снаружи вместо прямых импортов app.* (см. scraper_kit.contracts). process_houses_imv_batch зовётся через enrichment (EnrichmentJobs Protocol). """ _anchors = anchors if anchors is not None else EKB_ANCHORS counters = CitySweepCounters(anchors_total=len(_anchors)) all_touched_house_ids: set[int] = set() # Масштабируемый watchdog-таймаут для Avito anchor'а. В отличие от Yandex/Cian, # Avito detail-fetch идёт через browser service с page.goto timeout 60s. При # detail_top_n=20 detail-фаза может занять до N×60s — фиксированный ANCHOR_TIMEOUT_SEC # обрезает её раньше SERP, теряя счётчики. Масштабирование по detail_top_n + houses # гарантирует, что SERP+save всегда успевает до watchdog. _avito_anchor_timeout = ( float(ANCHOR_TIMEOUT_SEC) + detail_top_n * _AVITO_PER_DETAIL_S + (_AVITO_HOUSES_BUDGET_S if enrich_houses else 0.0) ) logger.info( "city-sweep run_id=%d: avito anchor_timeout=%.0fs " "(base=%d + detail_top_n=%d×%.0fs + houses=%.0fs)", run_id, _avito_anchor_timeout, ANCHOR_TIMEOUT_SEC, detail_top_n, _AVITO_PER_DETAIL_S, _AVITO_HOUSES_BUDGET_S if enrich_houses else 0.0, ) browser_mode = config.scraper_fetch_mode == "browser" # #1790: общий бюджет IP-ротаций на весь sweep (разделяется по всем anchor'ам). sweep_rotations_done = 0 async with AsyncExitStack() as stack: session: AsyncSession | None = None shared_bf: BrowserFetcher | None = None if browser_mode: shared_bf = await stack.enter_async_context( BrowserFetcher( source="avito", endpoint=config.browser_http_endpoint, proxy_provider=proxy_provider, use_pool=config.use_proxy_pool_browser, ) ) else: session = await stack.enter_async_context( AsyncSession( impersonate="chrome120", timeout=25, headers=_CHROME_HEADERS, proxies=_avito_proxies(config), ) ) try: for idx, (lat, lon, name) in enumerate(_anchors, start=1): if runs.is_cancelled(db, run_id): logger.info( "city-sweep run_id=%d: cancelled at anchor #%d/%d (%s)", run_id, idx, len(_anchors), name, ) runs.update_heartbeat(db, run_id, counters.to_dict()) return counters elif shutdown_requested(): # Кооперативный SIGTERM-drain (#1182 Phase 3a): останавливаемся # на границе anchor'а и финализируем run как mark_done(partial) — # это НЕ user-cancel, поэтому НЕ mark_cancelled. Дальнейшие anchor'ы # и IMV-фаза не выполняются. logger.info( "city-sweep run_id=%d: SIGTERM-drain — stopping at anchor #%d/%d (%s)", run_id, idx, len(_anchors), name, ) runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_done(db, run_id, counters.to_dict()) return counters logger.info( "city-sweep run_id=%d anchor #%d/%d (%s, %.4f, %.4f)", run_id, idx, len(_anchors), name, lat, lon, ) # Capture loop variables in default args (B023): prevents stale binding # if the coroutine is scheduled after the loop variable changes. _a_lat, _a_lon, _a_name = lat, lon, name async def _avito_anchor_phases( _lat: float = _a_lat, _lon: float = _a_lon, _name: str = _a_name, ) -> None: """Все фазы одного Avito anchor'а (SERP+save → houses → detail). Обновляет sweep-level counters и all_touched_house_ids через nonlocal НЕМЕДЛЕННО после SERP+save (до detail-фазы). Это гарантирует, что TimeoutError из wait_for теряет только detail-счётчики, но не SERP. """ nonlocal all_touched_house_ids, sweep_rotations_done # ── Phase 1+2: SERP + save ───────────────────────────────── scraper = AvitoScraper(config) if browser_mode: scraper._browser = shared_bf # Shared-browser режим: _cffi=None → curl_cffi-fallback на # firewall пропускается (by design). Логируем для observability. logger.debug("avito pipeline-mode: shared browser only, no cffi fallback") else: scraper._cffi = session try: anchor_lots: list[ScrapedLot] = await scraper.fetch_around( _lat, _lon, radius_m, pages=pages_per_anchor, delay_override_sec=request_delay_sec, ) except (AvitoBlockedError, AvitoRateLimitedError): logger.error( "city-sweep run_id=%d: SERP BLOCKED at anchor (%s, %.4f, %.4f)", run_id, _name, _lat, _lon, ) raise counters.lots_fetched += len(anchor_lots) if anchor_lots: try: ins, upd = save_listings( db, anchor_lots, matcher=matcher, region_code=region_code ) counters.lots_inserted += ins counters.lots_updated += upd except Exception as save_exc: logger.exception( "city-sweep run_id=%d: save_listings failed at anchor %s: %s", run_id, _name, save_exc, ) counters.errors_count += 1 logger.info( "city-sweep run_id=%d anchor %s: SERP fetched=%d ins=%d upd=%d", run_id, _name, len(anchor_lots), counters.lots_inserted, counters.lots_updated, ) # SERP-счётчики зафиксированы в sweep-level counters. # Дальнейший TimeoutError из wait_for их не сотрёт. # ── Phase 3: group by house ──────────────────────────────── unique_house_paths: set[str] = set() for lot in anchor_lots: if lot.house_url: try: parsed = urlparse(lot.house_url) path = parsed.path if parsed.path else lot.house_url unique_house_paths.add(path) except Exception: continue if unique_house_paths: counters.unique_houses += len(unique_house_paths) # ── Phase 4: enrich houses ───────────────────────────────── touched_house_ids: set[int] = set() sweep_consecutive_house_blocks = 0 # #1790: per-anchor счётчик if enrich_houses and unique_house_paths: house_paths_list = list(unique_house_paths) h_idx = 0 while h_idx < len(house_paths_list): house_path = house_paths_list[h_idx] try: enrichment_hc = await fetch_house_catalog( house_path, cffi_session=session, browser_fetcher=shared_bf, ) hc = save_house_catalog_enrichment( db, enrichment_hc, matcher=matcher ) counters.houses_enriched += 1 sweep_consecutive_house_blocks = 0 hid = hc.get("house_id") if hid: touched_house_ids.add(int(hid)) if h_idx < len(house_paths_list) - 1: await asyncio.sleep(request_delay_sec) h_idx += 1 except (AvitoBlockedError, AvitoRateLimitedError) as be: # #1790: пробуем сменить IP перед re-raise. sweep_consecutive_house_blocks += 1 counters.houses_failed += 1 counters.errors_count += 1 logger.warning( "city-sweep run_id=%d: house BLOCKED at %s " "(consecutive=%d) — %s", run_id, house_path, sweep_consecutive_house_blocks, be, ) if ( sweep_consecutive_house_blocks >= _AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT and sweep_rotations_done < config.avito_proxy_max_rotations ): rotated, sweep_rotations_done = ( await _try_rotate_within_budget( config, reason=( f"city-sweep house " f"consecutive={sweep_consecutive_house_blocks}" ), rotations_done=sweep_rotations_done, ) ) if rotated: sweep_consecutive_house_blocks = 0 logger.info( "city-sweep run_id=%d: house ROTATE — " "IP changed, reset consecutive " "(rotation #%d/%d)", run_id, sweep_rotations_done, config.avito_proxy_max_rotations, ) # Не инкрементируем h_idx — retry текущего дома continue elif sweep_consecutive_house_blocks >= 3: logger.error( "city-sweep run_id=%d: house ABORT — " "%d consecutive blocks, propagating", run_id, sweep_consecutive_house_blocks, ) raise h_idx += 1 except Exception as he: logger.warning( "city-sweep run_id=%d: house_enrich failed %s: %s", run_id, house_path, he, ) counters.houses_failed += 1 counters.errors_count += 1 try: db.rollback() except Exception: pass h_idx += 1 # ── Phase 5: enrich detail (top-N) ──────────────────────── if detail_top_n > 0 and anchor_lots: priority_rows = ( db.execute( text(""" SELECT source_url FROM listings WHERE source = 'avito' AND source_url IS NOT NULL AND ( ( detail_enriched_at IS NULL AND price_rub > 0 AND ST_DWithin( geom::geography, ST_MakePoint(:lon, :lat)::geography, :radius ) ) OR ( detail_enriched_at IS NULL AND scraped_at > NOW() - INTERVAL '2 hours' ) ) ORDER BY scraped_at DESC NULLS LAST LIMIT :limit """), { "lat": _lat, "lon": _lon, "radius": radius_m * 2, "limit": detail_top_n, }, ) .mappings() .all() ) consecutive_blocks = 0 consecutive_timeouts = 0 detail_house_paths: set[str] = set() for d_idx, row in enumerate(priority_rows): source_url: str = row["source_url"] counters.detail_attempted += 1 try: item_url = ( urlparse(source_url).path if source_url.startswith("http") else source_url ) enrichment_detail = await fetch_detail( item_url, cffi_session=session, browser_fetcher=shared_bf, ) if save_detail_enrichment(db, enrichment_detail): counters.detail_enriched += 1 if enrichment_detail.house_catalog_url: hp = urlparse(enrichment_detail.house_catalog_url).path if hp: detail_house_paths.add(hp) consecutive_blocks = 0 consecutive_timeouts = 0 except (AvitoBlockedError, AvitoRateLimitedError) as be: consecutive_blocks += 1 counters.detail_failed += 1 counters.errors_count += 1 logger.warning( "city-sweep run_id=%d: detail BLOCKED #%d/%d " "(consecutive=%d): %s", run_id, d_idx + 1, len(priority_rows), consecutive_blocks, be, ) if consecutive_blocks >= 3: # #1790: попытка ротации перед abort'ом. rotated, sweep_rotations_done = ( await _try_rotate_within_budget( config, reason=( f"city-sweep detail " f"consecutive={consecutive_blocks}" ), rotations_done=sweep_rotations_done, ) ) if rotated: consecutive_blocks = 0 logger.info( "city-sweep run_id=%d: detail ROTATE — " "IP changed, reset consecutive " "(rotation #%d/%d)", run_id, sweep_rotations_done, config.avito_proxy_max_rotations, ) continue logger.error( "city-sweep run_id=%d: detail ABORT — " "%d consecutive blocks (IP rate-limited)", run_id, consecutive_blocks, ) raise except Exception as de: counters.detail_failed += 1 counters.errors_count += 1 consecutive_timeouts += 1 logger.warning( "city-sweep run_id=%d: detail failed #%d/%d " "(consecutive_timeouts=%d): %s", run_id, d_idx + 1, len(priority_rows), consecutive_timeouts, de, ) # (C) Fail-fast: браузер явно деградирует — прерываем # detail-фазу, не тратим время на оставшиеся запросы. if consecutive_timeouts >= _AVITO_DETAIL_CONSECUTIVE_TIMEOUT_ABORT: logger.error( "city-sweep run_id=%d: detail ABORT — " "%d consecutive timeouts/errors, browser degraded", run_id, consecutive_timeouts, ) break try: db.rollback() except Exception: pass if d_idx < len(priority_rows) - 1: jitter = random.uniform(0.8, 1.2) await asyncio.sleep(request_delay_sec * jitter) # ── Phase 5b: houses найденные через detail-страницы ── new_house_paths = detail_house_paths - unique_house_paths if enrich_houses and new_house_paths: counters.unique_houses += len(new_house_paths) nh_list = list(new_house_paths) nh_idx = 0 sweep_consecutive_detail_house_blocks = 0 # #1790 while nh_idx < len(nh_list): house_path = nh_list[nh_idx] try: enrichment_hc = await fetch_house_catalog( house_path, cffi_session=session, browser_fetcher=shared_bf, ) hc = save_house_catalog_enrichment( db, enrichment_hc, matcher=matcher ) counters.houses_enriched += 1 sweep_consecutive_detail_house_blocks = 0 hid = hc.get("house_id") if hid: touched_house_ids.add(int(hid)) if nh_idx < len(nh_list) - 1: await asyncio.sleep(request_delay_sec) nh_idx += 1 except (AvitoBlockedError, AvitoRateLimitedError) as beh: # #1790: ротация перед re-raise. sweep_consecutive_detail_house_blocks += 1 counters.houses_failed += 1 counters.errors_count += 1 logger.warning( "city-sweep run_id=%d: house(detail) BLOCKED %s" " (consecutive=%d) — %s", run_id, house_path, sweep_consecutive_detail_house_blocks, beh, ) if ( sweep_consecutive_detail_house_blocks >= 3 and sweep_rotations_done < config.avito_proxy_max_rotations ): rotated, sweep_rotations_done = ( await _try_rotate_within_budget( config, reason=( "city-sweep house(detail) " f"consecutive=" f"{sweep_consecutive_detail_house_blocks}" ), rotations_done=sweep_rotations_done, ) ) if rotated: sweep_consecutive_detail_house_blocks = 0 logger.info( "city-sweep run_id=%d: house(detail) " "ROTATE — IP changed (rotation #%d/%d)", run_id, sweep_rotations_done, config.avito_proxy_max_rotations, ) continue elif sweep_consecutive_detail_house_blocks >= 3: logger.error( "city-sweep run_id=%d: house(detail) ABORT" " — %d consecutive blocks, propagating", run_id, sweep_consecutive_detail_house_blocks, ) raise nh_idx += 1 except Exception as he: logger.warning( "city-sweep run_id=%d: house(detail) failed %s: %s", run_id, house_path, he, ) counters.houses_failed += 1 counters.errors_count += 1 try: db.rollback() except Exception: pass nh_idx += 1 all_touched_house_ids.update(touched_house_ids) logger.info( "city-sweep run_id=%d anchor %s done: " "lots=%d houses=%d/%d detail=%d/%d touched=%d errors=%d", run_id, _name, len(anchor_lots), counters.houses_enriched, counters.unique_houses, counters.detail_enriched, counters.detail_attempted, len(touched_house_ids), counters.errors_count, ) try: await asyncio.wait_for(_avito_anchor_phases(), timeout=_avito_anchor_timeout) except TimeoutError: logger.warning( "city-sweep run_id=%d: anchor #%d/%d (%s, %.4f, %.4f) " "timed out after %.0fs — SERP counters preserved, " "detail phase incomplete", run_id, idx, len(_anchors), name, lat, lon, _avito_anchor_timeout, ) counters.errors_count += 1 except (AvitoBlockedError, AvitoRateLimitedError) as e: logger.error( "city-sweep run_id=%d ABORT at anchor #%d/%d (%s) — blocked: %s", run_id, idx, len(_anchors), name, e, ) counters.errors_count += 1 counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) # #1950: если SERP уже собрал лоты и заблокировало только detail/houses, # ставим 'done' (не 'banned') — partial intake сохранён. # За флагом avito_serp_ok_not_banned (default True). _serp_ok = (counters.lots_inserted + counters.lots_updated) > 0 if config.avito_serp_ok_not_banned and _serp_ok: _note = ( f"detail enrichment aborted ({e}) но SERP intake завершён: " f"{counters.lots_inserted} ins {counters.lots_updated} upd" ) logger.warning( "city-sweep run_id=%d: SERP OK (ins=%d upd=%d), " "detail blocked — DONE not banned. %s", run_id, counters.lots_inserted, counters.lots_updated, _note, ) runs.mark_done( db, run_id, {**counters.to_dict(), "enrichment_abort_note": _note}, # type: ignore[arg-type] ) else: runs.mark_banned(db, run_id, str(e), counters.to_dict()) return counters except Exception: logger.exception("city-sweep run_id=%d: anchor %s failed", run_id, name) counters.errors_count += 1 counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) # ── IMV-фаза: финальный обход тронутых домов ────────── if enrich_imv and all_touched_house_ids: if runs.is_cancelled(db, run_id): logger.info( "city-sweep run_id=%d: cancelled before IMV phase (%d houses skipped)", run_id, len(all_touched_house_ids), ) runs.mark_done(db, run_id, counters.to_dict()) return counters elif shutdown_requested(): # SIGTERM-drain до IMV-фазы: финализируем без дорогой IMV-оценки. logger.info( "city-sweep run_id=%d: SIGTERM-drain — skipping IMV phase " "(%d houses skipped)", run_id, len(all_touched_house_ids), ) runs.mark_done(db, run_id, counters.to_dict()) return counters logger.info( "city-sweep run_id=%d: IMV phase — %d touched houses", run_id, len(all_touched_house_ids), ) try: # Периодический heartbeat внутри длинной IMV-фазы (#1363): # без него десятки/сотни домов × request_delay_sec проходят # без update_heartbeat, и reap_zombies помечает живой run # 'zombie' → последующий mark_done становится no-op (дубль-sweep). def _imv_heartbeat() -> None: runs.update_heartbeat(db, run_id, counters.to_dict()) imv_result = await enrichment.process_houses_imv_batch( db, all_touched_house_ids, request_delay_sec=request_delay_sec, heartbeat=_imv_heartbeat, ) counters.imv_attempted += imv_result.checked counters.imv_enriched += imv_result.saved counters.imv_failed += imv_result.errors counters.errors_count += imv_result.errors runs.update_heartbeat(db, run_id, counters.to_dict()) logger.info( "city-sweep run_id=%d: IMV phase done — attempted=%d enriched=%d failed=%d", run_id, imv_result.checked, imv_result.saved, imv_result.errors, ) except Exception as exc: logger.error( "city-sweep run_id=%d: IMV phase crashed — %r (sweep still done)", run_id, exc, ) counters.errors_count += 1 # #1521: IMV-фаза могла оставить Session в aborted-state # (необёрнутый eligibility SELECT / _mark_status на отравленной # сессии). update_heartbeat и mark_done — единственные финализаторы # БЕЗ defensive rollback, поэтому сбрасываем txn-state сами. Иначе # PendingRollbackError пролетает мимо внутреннего try во внешний # except (mark_failed) → успешный sweep ложно помечается 'failed'. try: db.rollback() except Exception: pass runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_done(db, run_id, counters.to_dict()) logger.info( "city-sweep run_id=%d done: anchors=%d/%d lots=%d (ins=%d/upd=%d) " "houses=%d/%d detail=%d/%d imv=%d/%d errors=%d", run_id, counters.anchors_done, counters.anchors_total, counters.lots_fetched, counters.lots_inserted, counters.lots_updated, counters.houses_enriched, counters.unique_houses, counters.detail_enriched, counters.detail_attempted, counters.imv_enriched, counters.imv_attempted, counters.errors_count, ) return counters except Exception as exc: logger.exception("city-sweep run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) raise # ── Avito newbuilding (novostroyka) citywide sweep ──────────────── @dataclass class NewbuildingSweepCounters: """Aggregate counters для Avito novostroyka citywide sweep run.""" lots_fetched: int = 0 lots_inserted: int = 0 lots_updated: int = 0 errors_count: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} async def run_avito_newbuilding_sweep( db: Session, *, run_id: int, config: ScraperConfig, matcher: HouseMatcher, shutdown_requested: Callable[[], bool] = lambda: False, proxy_provider: ProxyProvider | None = None, pages: int = 20, request_delay_sec: float = 7.0, region_code: int = DEFAULT_REGION_CODE, ) -> NewbuildingSweepCounters: """Citywide-обход ЕКБ-выборки только новостроек (novostroyka-filter) → save. В отличие от run_avito_city_sweep (anchor-based fetch_around + houses/detail/IMV), это простой citywide paginated-обход через AvitoScraper.fetch_newbuildings: dedicated novostroyka-SERP отдаёт 100% new-build карточки (listing_segment= "novostroyki", newbuilding_id/newbuilding_url заполнены парсером). Сохраняем через стандартный save_listings, ведём SERP-счётчики через runs. - Single shared AsyncSession/BrowserFetcher на весь sweep (один TLS fingerprint) - AvitoBlockedError/AvitoRateLimitedError → mark_banned (status='banned') - is_cancelled-чек до старта SERP-фазы Инжекция (#2135 F2): config/matcher/shutdown_requested приходят снаружи вместо прямых импортов app.* (см. scraper_kit.contracts). """ counters = NewbuildingSweepCounters() browser_mode = config.scraper_fetch_mode == "browser" async with AsyncExitStack() as stack: session: AsyncSession | None = None shared_bf: BrowserFetcher | None = None if browser_mode: shared_bf = await stack.enter_async_context( BrowserFetcher( source="avito", endpoint=config.browser_http_endpoint, proxy_provider=proxy_provider, use_pool=config.use_proxy_pool_browser, ) ) else: session = await stack.enter_async_context( AsyncSession( impersonate="chrome120", timeout=25, headers=_CHROME_HEADERS, proxies=_avito_proxies(config), ) ) try: if runs.is_cancelled(db, run_id): logger.info("nb-sweep run_id=%d: cancelled before SERP phase", run_id) runs.mark_done(db, run_id, counters.to_dict()) return counters elif shutdown_requested(): logger.info( "nb-sweep run_id=%d: SIGTERM-drain — stopping before SERP phase", run_id ) runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_done(db, run_id, counters.to_dict()) return counters scraper = AvitoScraper(config) if browser_mode: scraper._browser = shared_bf # Shared-browser режим: _cffi=None → curl_cffi-fallback на firewall # пропускается (by design). Логируем для observability. logger.debug("avito pipeline-mode: shared browser only, no cffi fallback") else: scraper._cffi = session try: lots: list[ScrapedLot] = await scraper.fetch_newbuildings( pages=pages, delay_override_sec=request_delay_sec, ) except (AvitoBlockedError, AvitoRateLimitedError) as e: logger.error("nb-sweep run_id=%d: SERP BLOCKED — %s", run_id, e) runs.mark_banned(db, run_id, str(e), counters.to_dict()) return counters counters.lots_fetched += len(lots) if lots: try: ins, upd = save_listings(db, lots, matcher=matcher, region_code=region_code) counters.lots_inserted += ins counters.lots_updated += upd except Exception as save_exc: logger.exception( "nb-sweep run_id=%d: save_listings failed: %s", run_id, save_exc ) counters.errors_count += 1 try: db.rollback() except Exception: pass runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_done(db, run_id, counters.to_dict()) logger.info( "nb-sweep run_id=%d done: lots=%d (ins=%d/upd=%d) errors=%d", run_id, counters.lots_fetched, counters.lots_inserted, counters.lots_updated, counters.errors_count, ) return counters except Exception as exc: logger.exception("nb-sweep run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) raise # ── Yandex city sweep ─────────────────────────────────────────── # Константы для расчёта watchdog-таймаута combos-режима Yandex sweep. # В combos-режиме единственный "anchor" делает num_combos × max_pages fetch'ей; # каждый fetch занимает request_delay_sec + сетевой overhead. ADDRESS_ENRICH_BUDGET_S # добавляет буфер на address-enrich фазу (per-listing HTTP + sleep × N листингов). _YANDEX_COMBOS_PER_FETCH_S: float = 12.0 # network + parse margin (секунды на страницу) _YANDEX_ADDRESS_ENRICH_BUDGET_S: float = 300.0 # бюджет на address-enrich фазу @dataclass class YandexCitySweepCounters: """Aggregate counters для Yandex.Недвижимость city sweep run. Фазы: SERP (fetch_around_multi_room) → save_listings инкрементально → address-enrich. address_* — результат T10-обогащения (detail → полный адрес с номером дома). combos_skipped — сколько (room×price) combo пропущено из-за таймаутов gate-API. """ anchors_total: int = 0 anchors_done: int = 0 lots_fetched: int = 0 lots_inserted: int = 0 lots_updated: int = 0 price_history_rows: int = 0 combos_skipped: int = 0 address_attempted: int = 0 address_enriched: int = 0 address_failed: int = 0 errors_count: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} # Defensive abort threshold: у YandexRealtyScraper нет block/rate-exception классов # (fetch_around глотает ошибки и возвращает []). Чтобы не молотить Yandex впустую при # системном блоке/бане IP, прерываем sweep после N подряд anchor'ов, упавших с # исключением. NB: пустой результат (0 lots) — это НЕ ошибка, счётчик не растёт. YANDEX_SWEEP_MAX_CONSECUTIVE_FAILURES = 3 async def run_yandex_city_sweep( db: Session, *, run_id: int, config: ScraperConfig, matcher: HouseMatcher, enrichment: EnrichmentJobs, shutdown_requested: Callable[[], bool] = lambda: False, proxy_provider: ProxyProvider | None = None, anchors: list[tuple[float, float, str]] | None = None, radius_m: int = 25000, pages_per_anchor: int = 2, request_delay_sec: float | None = None, enrich_address: bool = True, rooms_list: list[str] | None = None, price_ranges: list[tuple[int | None, int | None]] | None = None, segments: list[str] | None = None, region_code: int = DEFAULT_REGION_CODE, ) -> YandexCitySweepCounters: """Yandex.Недвижимость city sweep: rooms × price combos от центра ЕКБ → save → address-enrich. Вместо 5 географических anchor'ов итерирует ВСЕ combinations (room × price_range) из одной центральной точки ЕКБ (56.8400, 60.6050). radius_m=25000 охватывает весь город; price/room segmentation обходит Yandex SERP ~575-card cap без geo-шрапнели. Фазы (единственный центральный якорь): 1. SERP: YandexRealtyScraper.fetch_around_multi_room(...) — combos mode (on_combo save). 2. ADDRESS-ENRICH (если enrich_address=True): для yandex-листингов без номера дома запрашиваем detail-страницу и извлекаем полный адрес из <title>. Инжекция (#2135 F2): config/matcher/shutdown_requested + продуктовые хелперы (record_yandex_price_history / extract_address_from_title / address_has_house_number) через enrichment (EnrichmentJobs) вместо прямых импортов app.services.yandex_*. Возвращает YandexCitySweepCounters. mark_done вызывается ВСЕГДА (кроме cancel/ban/fatal). """ # EKB center — единственный anchor для citywide combos sweep. # anchors-параметр принимается для backward-compat тестов. if anchors is not None: _anchors: list[tuple[float, float, str]] = anchors else: _anchors = [(56.8400, 60.6050, "Центр (combos)")] _rooms_list = rooms_list or list(ROOM_PATH.keys()) _price_ranges = price_ranges or DEFAULT_PRICE_RANGES # newFlat-сегменты: оба прохода (vtorichka + novostroyki) по одним combos. _segments = segments or ["NO", "YES"] counters = YandexCitySweepCounters(anchors_total=len(_anchors)) inter_anchor_delay = request_delay_sec if request_delay_sec is not None else 7.0 enrich_delay = request_delay_sec if request_delay_sec is not None else 3.0 _resolved_delay = request_delay_sec if request_delay_sec is not None else 9.0 consecutive_failures = 0 yandex_rotations_done = 0 # #1848: бюджет IP-ротаций на весь sweep # Вычисляем watchdog-таймаут для combos-режима (центр, anchors=None). _num_combos = len(_rooms_list) * len(_price_ranges) if anchors is None and _num_combos > 0: _sweep_timeout = max( ANCHOR_TIMEOUT_SEC, len(_segments) * _num_combos * pages_per_anchor * (_resolved_delay + _YANDEX_COMBOS_PER_FETCH_S) + _YANDEX_ADDRESS_ENRICH_BUDGET_S, ) else: # Explicit anchors (backward-compat тесты или ручной override) — legacy timeout. _sweep_timeout = ANCHOR_TIMEOUT_SEC # Прокси для address-enrich curl_cffi сессии (зеркало yandex_address_backfill). _proxy_url = config.scraper_proxy_url _proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None try: for idx, (lat, lon, name) in enumerate(_anchors, start=1): # ── Cooperative cancel ─────────────────────────────────────────── if runs.is_cancelled(db, run_id): logger.info( "yandex-sweep run_id=%d: cancelled at anchor #%d/%d (%s)", run_id, idx, len(_anchors), name, ) runs.update_heartbeat(db, run_id, counters.to_dict()) return counters elif shutdown_requested(): # SIGTERM-drain: чистый стоп на границе anchor'а → mark_done(partial). logger.info( "yandex-sweep run_id=%d: SIGTERM-drain — stopping at anchor #%d/%d (%s)", run_id, idx, len(_anchors), name, ) runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_done(db, run_id, counters.to_dict()) return counters logger.info( "yandex-sweep run_id=%d center-combos #%d/%d (%s, %.4f, %.4f) " "radius_m=%d rooms=%d price_ranges=%d", run_id, idx, len(_anchors), name, lat, lon, radius_m, len(_rooms_list), len(_price_ranges), ) # ── Per-anchor phases под watchdog-таймаутом ──────────────────── anchor_lots: list[ScrapedLot] = [] _anchor_timed_out = False _lat, _lon, _name = lat, lon, name _combo_lots_accumulator: list[ScrapedLot] = anchor_lots async def _yandex_anchor_phases( _a_lat: float = _lat, _a_lon: float = _lon, _a_name: str = _name, _al: list[ScrapedLot] = _combo_lots_accumulator, ) -> None: """Все фазы одного anchor'а (SERP + address-enrich).""" nonlocal consecutive_failures, yandex_rotations_done # ── Phase 1+2: SERP + инкрементальный save per-combo ────── def _on_combo( combo_label: str, new_lots: list[ScrapedLot], _accumulator: list[ScrapedLot] = _al, ) -> None: """Callback: вызывается scraper'ом после каждого combo с новыми лотами.""" _accumulator.extend(new_lots) counters.lots_fetched += len(new_lots) try: ins, upd = save_listings( db, new_lots, matcher=matcher, region_code=region_code, run_id=run_id ) counters.lots_inserted += ins counters.lots_updated += upd except Exception as save_exc: counters.errors_count += 1 logger.warning( "yandex-sweep run_id=%d combo %s: save_listings failed: %s", run_id, combo_label, save_exc, ) try: db.rollback() except Exception: pass runs.update_heartbeat(db, run_id, counters.to_dict()) logger.debug( "yandex-sweep run_id=%d combo %s: saved %d lots " "(ins=%d upd=%d total_fetched=%d)", run_id, combo_label, len(new_lots), counters.lots_inserted, counters.lots_updated, counters.lots_fetched, ) async with YandexRealtyScraper(config, proxy_provider=proxy_provider) as scraper: await scraper.fetch_around_multi_room( _a_lat, _a_lon, radius_m, max_pages=pages_per_anchor, rooms_list=_rooms_list, price_ranges=_price_ranges, segments=_segments, on_combo=_on_combo, ) # _al накоплен on_combo; save уже вызван per-combo. # Price-history из gate price.previous/trend (graceful — провал # истории НЕ должен ронять сейв листингов). if _al: try: counters.price_history_rows += enrichment.record_yandex_price_history( db, _al ) except Exception as ph_exc: logger.warning( "yandex-sweep run_id=%d: price-history failed: %s", run_id, ph_exc, ) try: db.rollback() except Exception: pass consecutive_failures = 0 logger.info( "yandex-sweep run_id=%d center-combos %s: fetched=%d ins=%d upd=%d", run_id, _a_name, counters.lots_fetched, counters.lots_inserted, counters.lots_updated, ) # ── Phase 3: address-enrich ──────────────────────────────── if not (enrich_address and _al): return source_ids = [lot.source_id for lot in _al if lot.source_id] if not source_ids: return id_rows = ( db.execute( text(""" SELECT id, source_id, address, source_url FROM listings WHERE source = 'yandex' AND source_id = ANY(CAST(:ids AS text[])) AND address IS NOT NULL AND NOT (address ~ ',\\s*\\d+') AND source_url IS NOT NULL """), {"ids": source_ids}, ) .mappings() .all() ) items = [dict(r) for r in id_rows] counters.address_attempted += len(items) if items: # #1848: rotation tracking для Yandex address-enrich (stuck-proxy) _yandex_enrich_consec_fail = 0 _yandex_enrich_abort = 3 # abort address-enrich после N подряд async with AsyncSession( impersonate="chrome120", timeout=30.0, proxies=_proxies, headers={"Accept-Language": "ru-RU,ru;q=0.9,en;q=0.8"}, ) as enrich_session: eidx = 0 while eidx < len(items): item = items[eidx] lid: int = item["id"] old_addr: str = item["address"] url: str = item["source_url"] try: resp = await enrich_session.get(url, allow_redirects=True) if resp.status_code not in {200}: logger.warning( "yandex-sweep run_id=%d address-enrich:" " HTTP %d listing_id=%d", run_id, resp.status_code, lid, ) counters.address_failed += 1 # 503/429 = прокси заблокирован → считаем как failure if resp.status_code in {429, 503}: _yandex_enrich_consec_fail += 1 else: _yandex_enrich_consec_fail = 0 else: new_addr = enrichment.extract_address_from_title(resp.text) if ( new_addr is None or new_addr == old_addr or not enrichment.address_has_house_number(new_addr) ): pass else: try: db.execute( text(""" UPDATE listings SET address = :addr, geocode_tried_at = NULL WHERE id = CAST(:id AS bigint) """), {"addr": new_addr, "id": lid}, ) db.commit() counters.address_enriched += 1 logger.info( "yandex-sweep run_id=%d: " "address updated listing_id=%d " "%r -> %r", run_id, lid, old_addr, new_addr, ) except Exception as upd_exc: counters.address_failed += 1 counters.errors_count += 1 logger.warning( "yandex-sweep run_id=%d: " "address DB update failed " "listing_id=%d: %s", run_id, lid, upd_exc, ) try: db.rollback() except Exception: pass _yandex_enrich_consec_fail = 0 except Exception as fetch_exc: counters.address_failed += 1 _yandex_enrich_consec_fail += 1 logger.warning( "yandex-sweep run_id=%d: address-enrich " "fetch failed listing_id=%d " "(consecutive=%d): %s", run_id, lid, _yandex_enrich_consec_fail, fetch_exc, ) # Ротация Yandex-прокси при stuck/ban (#1848) if _yandex_enrich_consec_fail >= _yandex_enrich_abort: rotated, yandex_rotations_done = await _try_rotate_within_budget( config, reason=( f"yandex enrich consecutive={_yandex_enrich_consec_fail}" ), rotations_done=yandex_rotations_done, source="yandex", ) if rotated: _yandex_enrich_consec_fail = 0 logger.info( "yandex-sweep run_id=%d: enrich ROTATE — " "IP changed, reset consecutive " "(rotation #%d/%d)", run_id, yandex_rotations_done, config.yandex_proxy_max_rotations, ) # retry текущего item после ротации continue logger.error( "yandex-sweep run_id=%d: address-enrich ABORT — " "%d consecutive failures (proxy likely stuck/banned)", run_id, _yandex_enrich_consec_fail, ) break if eidx < len(items) - 1: await asyncio.sleep(enrich_delay) eidx += 1 logger.info( "yandex-sweep run_id=%d anchor %s: address enrich=%d/%d failed=%d", run_id, _a_name, counters.address_enriched, counters.address_attempted, counters.address_failed, ) try: await asyncio.wait_for(_yandex_anchor_phases(), timeout=_sweep_timeout) except TimeoutError: logger.warning( "yandex-sweep run_id=%d: anchor #%d/%d (%s, %.4f, %.4f) " "timed out after %.0fs — skipping", run_id, idx, len(_anchors), name, lat, lon, _sweep_timeout, ) counters.errors_count += 1 consecutive_failures += 1 _anchor_timed_out = True if consecutive_failures >= YANDEX_SWEEP_MAX_CONSECUTIVE_FAILURES: logger.error( "yandex-sweep run_id=%d ABORT at anchor #%d/%d (%s) — " "%d consecutive failures (last: timeout): " "IP likely blocked or proxy hung", run_id, idx, len(_anchors), name, consecutive_failures, ) counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_banned( db, run_id, f"yandex sweep aborted: {consecutive_failures} consecutive " f"anchor failures (last: anchor timeout {_sweep_timeout:.0f}s)", counters.to_dict(), ) return counters counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) if idx < len(_anchors): await asyncio.sleep(inter_anchor_delay) continue except Exception as e: # YandexRealtyScraper не propagate'ит block/rate exceptions — ловим broad. logger.exception("yandex-sweep run_id=%d: anchor %s SERP failed", run_id, name) counters.errors_count += 1 consecutive_failures += 1 if consecutive_failures >= YANDEX_SWEEP_MAX_CONSECUTIVE_FAILURES: logger.error( "yandex-sweep run_id=%d ABORT at anchor #%d/%d (%s) — " "%d consecutive failures (IP likely blocked): %s", run_id, idx, len(_anchors), name, consecutive_failures, e, ) counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_banned( db, run_id, f"yandex sweep aborted: {consecutive_failures} consecutive " f"anchor failures (last: {e})", counters.to_dict(), ) return counters counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) # Пауза перед следующим anchor'ом даже при ошибке. if idx < len(_anchors): await asyncio.sleep(inter_anchor_delay) continue counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) # Доп. пауза между anchor'ами (поверх per-page sleep внутри scraper'а). if idx < len(_anchors): await asyncio.sleep(inter_anchor_delay) runs.mark_done(db, run_id, counters.to_dict()) logger.info( "yandex-sweep run_id=%d done: anchors=%d/%d lots=%d (ins=%d/upd=%d) " "address=%d/%d failed=%d errors=%d", run_id, counters.anchors_done, counters.anchors_total, counters.lots_fetched, counters.lots_inserted, counters.lots_updated, counters.address_enriched, counters.address_attempted, counters.address_failed, counters.errors_count, ) return counters except Exception as exc: logger.exception("yandex-sweep run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) raise # ── Cian city sweep (#860) ────────────────────────────────────────────────── @dataclass class CianCitySweepCounters: """Aggregate counters для Cian city sweep run. Поля симметричны CitySweepCounters, но без imv_* (у Cian нет IMV-фазы). detail_* — обогащение через cian_detail (вкл. price-history). houses_* — обогащение newbuilding/ЖК через cian_newbuilding. """ anchors_total: int = 0 anchors_done: int = 0 lots_fetched: int = 0 lots_dropped_secondary: int = 0 # вторичка, отброшенная при newbuilding_only=True lots_inserted: int = 0 lots_updated: int = 0 detail_attempted: int = 0 detail_enriched: int = 0 detail_failed: int = 0 houses_attempted: int = 0 houses_enriched: int = 0 houses_failed: int = 0 errors_count: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} # Defensive abort threshold: если N подряд anchor'ов упали с исключением # (не просто вернули пустой результат) — прерываем sweep. CIAN_SWEEP_MAX_CONSECUTIVE_FAILURES = 3 async def run_cian_city_sweep( db: Session, *, run_id: int, config: ScraperConfig, matcher: HouseMatcher, shutdown_requested: Callable[[], bool] = lambda: False, proxy_provider: ProxyProvider | None = None, anchors: list[tuple[float, float, str]] | None = None, radius_m: int = 1500, pages_per_anchor: int = 3, request_delay_sec: float = 5.0, detail_top_n: int = 10, enrich_houses: bool = True, newbuilding_only: bool = True, region_code: int = DEFAULT_REGION_CODE, ) -> CianCitySweepCounters: """Cian newbuilding city sweep: SERP → detail(+price-history) → newbuilding/houses. Симметрично run_avito_city_sweep. Использует EKB_ANCHORS (общий список). newbuilding_only (default True): в SERP-фазе сохраняются только лоты с listing_segment == "novostroyki" — вторичку отбрасываем, ею авторитетно владеет run_cian_full_load (exhaustive региональный сбор). Инжекция (#2135 F2): config/matcher/shutdown_requested приходят снаружи вместо прямых импортов app.* (см. scraper_kit.contracts). Cooperative cancel: runs.is_cancelled(db, run_id) перед каждым anchor. mark_done вызывается ВСЕГДА (finally outer). """ _anchors = anchors if anchors is not None else EKB_ANCHORS counters = CianCitySweepCounters(anchors_total=len(_anchors)) consecutive_failures = 0 cian_rotations_done = 0 # #1848: бюджет IP-ротаций на весь sweep # #2160: масштабируемый watchdog-таймаут для cian anchor'а. При proxy-pool browser каждый # SERP-фетч 13-45s (camoufox relaunch), фиксированный ANCHOR_TIMEOUT_SEC=240 гильотинит # якорь в SERP-фазе. См. _cian_anchor_timeout_s + константы выше. _cian_anchor_timeout = _cian_anchor_timeout_s( pages_per_anchor, request_delay_sec, detail_top_n, enrich_houses ) logger.info( "cian-sweep run_id=%d: anchor_timeout=%.0fs " "(base=%d, serp=%d×%d×(%.0f+%.0f)s + detail_top_n=%d×%.0fs + houses=%.0fs)", run_id, _cian_anchor_timeout, ANCHOR_TIMEOUT_SEC, _CIAN_ROOM_BUCKETS, pages_per_anchor, request_delay_sec, _CIAN_PER_SERP_FETCH_S, detail_top_n, _CIAN_PER_DETAIL_S, _CIAN_HOUSES_BUDGET_S if enrich_houses else 0.0, ) try: for idx, (lat, lon, name) in enumerate(_anchors, start=1): # Cooperative cancel перед каждым anchor if runs.is_cancelled(db, run_id): logger.info( "cian-sweep run_id=%d: cancelled at anchor #%d/%d (%s)", run_id, idx, len(_anchors), name, ) runs.update_heartbeat(db, run_id, counters.to_dict()) return counters elif shutdown_requested(): # SIGTERM-drain: чистый стоп на границе anchor'а → mark_done(partial). logger.info( "cian-sweep run_id=%d: SIGTERM-drain — stopping at anchor #%d/%d (%s)", run_id, idx, len(_anchors), name, ) runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_done(db, run_id, counters.to_dict()) return counters logger.info( "cian-sweep run_id=%d anchor #%d/%d (%s, %.4f, %.4f)", run_id, idx, len(_anchors), name, lat, lon, ) # ── Per-anchor phases под watchdog-таймаутом ──────────────────── anchor_lots: list[ScrapedLot] = [] _c_lat, _c_lon, _c_name = lat, lon, name async def _cian_anchor_phases( _a_lat: float = _c_lat, _a_lon: float = _c_lon, _a_name: str = _c_name, ) -> None: """Все фазы одного cian anchor'а (SERP + detail + houses).""" nonlocal anchor_lots, consecutive_failures, cian_rotations_done # ── Phase 1+2: SERP + save ───────────────────────────────── async with CianScraper(config, proxy_provider=proxy_provider) as scraper: anchor_lots = await scraper.fetch_around_multi_room( _a_lat, _a_lon, radius_m, pages=pages_per_anchor, ) counters.lots_fetched += len(anchor_lots) # ── Фильтр вторички (newbuilding_only) ───────────────────── if newbuilding_only: _before = len(anchor_lots) anchor_lots = [ lot for lot in anchor_lots if lot.listing_segment == "novostroyki" ] counters.lots_dropped_secondary += _before - len(anchor_lots) if anchor_lots: inserted, updated = save_listings( db, anchor_lots, matcher=matcher, region_code=region_code, run_id=run_id ) counters.lots_inserted += inserted counters.lots_updated += updated consecutive_failures = 0 logger.info( "cian-sweep run_id=%d anchor %s: SERP fetched=%d nb_kept=%d " "dropped_secondary=%d ins=%d upd=%d", run_id, _a_name, counters.lots_fetched, len(anchor_lots), counters.lots_dropped_secondary, counters.lots_inserted, counters.lots_updated, ) # ── Phase 3: detail enrichment (top-N) ──────────────────── if detail_top_n > 0: priority_rows = ( db.execute( text(""" SELECT id, source_url FROM listings WHERE source = 'cian' AND source_url IS NOT NULL AND detail_enriched_at IS NULL AND scraped_at > NOW() - INTERVAL '2 hours' ORDER BY scraped_at DESC NULLS LAST LIMIT :lim """), {"lim": detail_top_n}, ) .mappings() .all() ) # #1848: Cian bans/stuck-proxy tracking + rotation (зеркало avito) _cian_detail_consec_failures = 0 _cian_detail_abort = 3 # abort detail-фазы после N подряд didx = 0 while didx < len(priority_rows): row = priority_rows[didx] listing_id: int = row["id"] source_url: str = row["source_url"] counters.detail_attempted += 1 try: cian_enrichment = await cian_fetch_detail( source_url, config=config, proxy_provider=proxy_provider ) if cian_enrichment is not None: cian_save_detail_enrichment( db, listing_id, cian_enrichment, matcher=matcher ) counters.detail_enriched += 1 else: counters.detail_failed += 1 logger.warning( "cian-sweep run_id=%d: detail returned None for " "listing_id=%d url=%s", run_id, listing_id, source_url, ) _cian_detail_consec_failures = 0 except Exception as exc: counters.detail_failed += 1 counters.errors_count += 1 _cian_detail_consec_failures += 1 logger.warning( "cian-sweep run_id=%d: detail failed listing_id=%d " "(consecutive=%d): %s", run_id, listing_id, _cian_detail_consec_failures, exc, ) if _cian_detail_consec_failures >= _cian_detail_abort: # Попытка ротации Cian-прокси перед abort'ом (#1848) rotated, cian_rotations_done = await _try_rotate_within_budget( config, reason=( f"cian detail consecutive={_cian_detail_consec_failures}" ), rotations_done=cian_rotations_done, source="cian", ) if rotated: _cian_detail_consec_failures = 0 logger.info( "cian-sweep run_id=%d: detail ROTATE — " "IP changed, reset consecutive " "(rotation #%d/%d)", run_id, cian_rotations_done, config.cian_proxy_max_rotations, ) # retry текущего detail (не инкрементируем didx) continue logger.error( "cian-sweep run_id=%d: detail ABORT — " "%d consecutive failures (proxy likely banned/stuck)", run_id, _cian_detail_consec_failures, ) break if didx < len(priority_rows) - 1: jitter = random.uniform(0.8, 1.2) await asyncio.sleep(request_delay_sec * jitter) didx += 1 logger.info( "cian-sweep run_id=%d anchor %s: detail=%d/%d failed=%d", run_id, _a_name, counters.detail_enriched, counters.detail_attempted, counters.detail_failed, ) # ── Phase 4: newbuilding/houses enrichment ───────────────── if not (enrich_houses and anchor_lots): return cian_nb_ext_ids = { lot.house_ext_id for lot in anchor_lots if lot.house_source == "cian_newbuilding" and lot.house_ext_id } if not cian_nb_ext_ids: return nb_id_list = list(cian_nb_ext_ids) try: house_rows = ( db.execute( text(""" SELECT h.id, h.cian_zhk_url FROM houses h JOIN house_sources hs ON hs.house_id = h.id WHERE hs.ext_source = 'cian_newbuilding' AND hs.ext_id = ANY(CAST(:ids AS text[])) AND h.cian_zhk_url IS NOT NULL """), {"ids": nb_id_list}, ) .mappings() .all() ) except Exception as exc: logger.exception( "cian-sweep run_id=%d anchor %s: houses DB query failed — " "skipping houses phase (%d ids): %s", run_id, _a_name, len(nb_id_list), exc, ) counters.houses_failed += len(nb_id_list) try: db.rollback() except Exception: pass return for hidx, hrow in enumerate(house_rows): house_id_val: int = hrow["id"] zhk_url: str = hrow["cian_zhk_url"] counters.houses_attempted += 1 try: nb_enrichment = await fetch_newbuilding(zhk_url, config=config) if nb_enrichment is not None: save_newbuilding_enrichment(db, house_id_val, nb_enrichment) counters.houses_enriched += 1 else: counters.houses_failed += 1 logger.warning( "cian-sweep run_id=%d: newbuilding returned None " "house_id=%d url=%s", run_id, house_id_val, zhk_url, ) except Exception as exc: counters.houses_failed += 1 counters.errors_count += 1 logger.warning( "cian-sweep run_id=%d: houses failed house_id=%d: %s", run_id, house_id_val, exc, ) try: db.rollback() except Exception: pass if hidx < len(house_rows) - 1: await asyncio.sleep(request_delay_sec) logger.info( "cian-sweep run_id=%d anchor %s: houses=%d/%d failed=%d", run_id, _a_name, counters.houses_enriched, counters.houses_attempted, counters.houses_failed, ) try: await asyncio.wait_for(_cian_anchor_phases(), timeout=_cian_anchor_timeout) except TimeoutError: logger.warning( "cian-sweep run_id=%d: anchor #%d/%d (%s, %.4f, %.4f) " "timed out after %ds — skipping", run_id, idx, len(_anchors), name, lat, lon, _cian_anchor_timeout, ) counters.errors_count += 1 consecutive_failures += 1 if consecutive_failures >= CIAN_SWEEP_MAX_CONSECUTIVE_FAILURES: logger.error( "cian-sweep run_id=%d ABORT at anchor #%d/%d (%s) — " "%d consecutive failures (last: timeout): " "IP likely blocked or proxy hung", run_id, idx, len(_anchors), name, consecutive_failures, ) counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_banned( db, run_id, f"cian sweep aborted: {consecutive_failures} consecutive " f"anchor failures (last: anchor timeout {_cian_anchor_timeout:.0f}s)", counters.to_dict(), ) return counters counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) continue except Exception as e: logger.exception("cian-sweep run_id=%d: anchor %s SERP failed", run_id, name) counters.errors_count += 1 consecutive_failures += 1 if consecutive_failures >= CIAN_SWEEP_MAX_CONSECUTIVE_FAILURES: logger.error( "cian-sweep run_id=%d ABORT at anchor #%d/%d (%s) — " "%d consecutive SERP failures (IP likely blocked): %s", run_id, idx, len(_anchors), name, consecutive_failures, e, ) counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_banned( db, run_id, f"cian sweep aborted: {consecutive_failures} consecutive " f"anchor SERP failures (last: {e})", counters.to_dict(), ) return counters counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) continue counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) # Пауза между anchor'ами (поверх per-request sleep внутри scraper'а) if idx < len(_anchors): await asyncio.sleep(request_delay_sec) runs.mark_done(db, run_id, counters.to_dict()) logger.info( "cian-sweep run_id=%d done: anchors=%d/%d lots=%d (ins=%d/upd=%d) " "detail=%d/%d houses=%d/%d errors=%d", run_id, counters.anchors_done, counters.anchors_total, counters.lots_fetched, counters.lots_inserted, counters.lots_updated, counters.detail_enriched, counters.detail_attempted, counters.houses_enriched, counters.houses_attempted, counters.errors_count, ) return counters except Exception as exc: logger.exception("cian-sweep run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) raise # ── Cian exhaustive full load (region-wide, без anchor'ов) ───────────────── @dataclass class CianFullLoadCounters: """Aggregate counters для run_cian_full_load. Без anchor-полей — единственный региональный проход. detail_* — опциональное обогащение detail-страниц (если enrich_detail=True). """ unique_fetched: int = 0 saved_inserted: int = 0 saved_updated: int = 0 detail_attempted: int = 0 detail_enriched: int = 0 detail_failed: int = 0 errors_count: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} async def run_cian_full_load( db: Session, *, run_id: int, config: ScraperConfig, matcher: HouseMatcher, shutdown_requested: Callable[[], bool] = lambda: False, proxy_provider: ProxyProvider | None = None, price_cap_per_bucket: int = 1400, detail_top_n: int = 0, request_delay_sec: float = 4.0, enrich_detail: bool = False, concurrency: int = 5, resume_run_id: int | None = None, region_code: int = DEFAULT_REGION_CODE, ) -> CianFullLoadCounters: """Exhaustive региональный сбор Cian ЕКБ вторички (БЕЗ anchor'ов). Один проход с партиционированием по КОМНАТНОСТИ × ЦЕНЕ (адаптивное бинарное деление) охватывает весь ЕКБ. Инкрементальный save: on_bucket коммитит каждый leaf-бакет в БД сразу после сбора. Инжекция (#2135 F2): config/matcher/shutdown_requested приходят снаружи вместо прямых импортов app.* (см. scraper_kit.contracts). Cooperative cancel: runs.is_cancelled проверяется per-bucket. """ counters = CianFullLoadCounters() # ── Checkpoint/resume: читаем done_buckets из прошлого run ─────────────── skip_set: set[str] = set() if resume_run_id is not None: _prev_row = db.execute( text("SELECT counters FROM scrape_runs WHERE id = CAST(:rid AS bigint)"), {"rid": resume_run_id}, ).fetchone() if _prev_row is not None and _prev_row.counters: _prev_counters: dict = ( _prev_row.counters if isinstance(_prev_row.counters, dict) else {} ) skip_set = set(_prev_counters.get("done_buckets", [])) logger.info( "cian-full-load run_id=%d: resuming from run %d — %d buckets already done", run_id, resume_run_id, len(skip_set), ) done: set[str] = set(skip_set) # накапливаем завершённые бакеты этого прогона def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg] """Инкрементальный save после каждого leaf-бакета. Дописывает bucket_key в done.""" nonlocal done if runs.is_cancelled(db, run_id): logger.info("cian-full-load run_id=%d: cancel detected in on_bucket", run_id) raise RuntimeError("cancelled") elif shutdown_requested(): # SIGTERM-drain: тот же sentinel-механизм, что и user-cancel, но отдельный # маркер "shutdown" → handler финализирует mark_done(partial), не mark_failed. logger.info( "cian-full-load run_id=%d: SIGTERM-drain — stopping at bucket %s", run_id, bucket_key, ) raise RuntimeError("shutdown") if not lots: done.add(bucket_key) runs.update_heartbeat( db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)} ) return inserted, updated = save_listings( db, lots, matcher=matcher, region_code=region_code, run_id=run_id, skip_seen_today=config.scraper_skip_seen_today, ) # save_listings вызывает db.commit() внутри — данные в БД сразу counters.saved_inserted += inserted counters.saved_updated += updated counters.unique_fetched += len(lots) done.add(bucket_key) runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) logger.info( "cian-full-load run_id=%d: bucket %s saved ins=%d upd=%d total_unique=%d", run_id, bucket_key, inserted, updated, counters.unique_fetched, ) def _on_progress(unique_count: int) -> None: """Heartbeat per room-bucket (на случай пустых бакетов без on_bucket вызовов).""" runs.update_heartbeat(db, run_id, counters.to_dict()) # #1949: фоновый heartbeat — обновляет heartbeat_at каждые 60s независимо от on_bucket. async def _background_heartbeat() -> None: """Периодический heartbeat каждые 60s, не завязанный на прогресс бакетов.""" while True: await asyncio.sleep(60.0) try: runs.update_heartbeat(db, run_id, counters.to_dict()) except Exception: pass _hb_task: asyncio.Task[None] | None = None try: _hb_task = asyncio.create_task(_background_heartbeat()) async with CianScraper(config, proxy_provider=proxy_provider) as scraper: scraper.request_delay_sec = request_delay_sec # #1949: per-fetch hard-cancel для browser-fetch'ей внутри bucket-gather. _fetch_timeout = config.cian_full_load_per_fetch_timeout_s if _fetch_timeout > 0 and scraper._browser is not None: _orig_fetch = scraper._browser.fetch async def _timed_fetch( url: str, _f: Any = _orig_fetch, _t: float = _fetch_timeout, _rid: int = run_id, ) -> str: try: return await asyncio.wait_for(_f(url), timeout=_t) except TimeoutError: logger.warning( "cian-full-load run_id=%d: browser-fetch timed out " "after %.0fs — url=%s", _rid, _t, url, ) raise scraper._browser.fetch = _timed_fetch # type: ignore[method-assign] logger.info( "cian-full-load run_id=%d: per-fetch timeout %.0fs enabled", run_id, _fetch_timeout, ) await scraper.fetch_all_secondary( price_cap_per_bucket=price_cap_per_bucket, concurrency=concurrency, secondary_only=True, on_bucket=_on_bucket, on_progress=_on_progress, skip_buckets=skip_set if skip_set else None, ) logger.info( "cian-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d", run_id, counters.unique_fetched, counters.saved_inserted, counters.saved_updated, ) runs.update_heartbeat(db, run_id, counters.to_dict()) # ── Опциональный detail-энричмент ───────────────────────────────────── if enrich_detail and detail_top_n > 0: priority_rows = ( db.execute( text(""" SELECT id, source_url FROM listings WHERE source = 'cian' AND source_url IS NOT NULL AND detail_enriched_at IS NULL AND scraped_at > NOW() - INTERVAL '4 hours' ORDER BY scraped_at DESC NULLS LAST LIMIT :lim """), {"lim": detail_top_n}, ) .mappings() .all() ) for didx, row in enumerate(priority_rows): if runs.is_cancelled(db, run_id): logger.info("cian-full-load run_id=%d: cancelled during detail enrich", run_id) break elif shutdown_requested(): logger.info( "cian-full-load run_id=%d: SIGTERM-drain — stopping detail at #%d/%d", run_id, didx + 1, len(priority_rows), ) break listing_id: int = row["id"] source_url: str = row["source_url"] counters.detail_attempted += 1 try: cian_enrichment = await cian_fetch_detail( source_url, config=config, proxy_provider=proxy_provider ) if cian_enrichment is not None: cian_save_detail_enrichment( db, listing_id, cian_enrichment, matcher=matcher ) counters.detail_enriched += 1 else: counters.detail_failed += 1 except Exception as exc: counters.detail_failed += 1 counters.errors_count += 1 logger.warning( "cian-full-load run_id=%d: detail failed listing_id=%d: %s", run_id, listing_id, exc, ) if didx < len(priority_rows) - 1: await asyncio.sleep(request_delay_sec * random.uniform(0.8, 1.2)) logger.info( "cian-full-load run_id=%d: detail done — enriched=%d/%d failed=%d", run_id, counters.detail_enriched, counters.detail_attempted, counters.detail_failed, ) runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) logger.info( "cian-full-load run_id=%d done: unique=%d ins=%d upd=%d detail=%d/%d errors=%d", run_id, counters.unique_fetched, counters.saved_inserted, counters.saved_updated, counters.detail_enriched, counters.detail_attempted, counters.errors_count, ) return counters except RuntimeError as exc: # on_bucket кидает RuntimeError("cancelled") при cooperative cancel, # RuntimeError("shutdown") при кооперативном SIGTERM-drain (#1182 Phase 3a). if str(exc) == "cancelled": logger.info( "cian-full-load run_id=%d: cancelled — partial results unique=%d ins=%d upd=%d", run_id, counters.unique_fetched, counters.saved_inserted, counters.saved_updated, ) runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) return counters if str(exc) == "shutdown": logger.info( "cian-full-load run_id=%d: SIGTERM-drain — partial results unique=%d ins=%d upd=%d", run_id, counters.unique_fetched, counters.saved_inserted, counters.saved_updated, ) runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) return counters logger.exception("cian-full-load run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) raise except Exception as exc: logger.exception("cian-full-load run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) raise finally: # #1949: отменяем фоновый heartbeat при любом завершении (return / raise / cancel). if _hb_task is not None: _hb_task.cancel() try: await _hb_task except asyncio.CancelledError: pass # ── Yandex exhaustive full load (region-wide, без anchor'ов) ────────────────── @dataclass class YandexFullLoadCounters: """Aggregate counters для run_yandex_full_load. Без detail_* полей — Yandex full load не делает detail-обогащение. """ unique_fetched: int = 0 saved_inserted: int = 0 saved_updated: int = 0 price_history_rows: int = 0 errors_count: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} async def run_yandex_full_load( db: Session, *, run_id: int, config: ScraperConfig, matcher: HouseMatcher, enrichment: EnrichmentJobs, shutdown_requested: Callable[[], bool] = lambda: False, proxy_provider: ProxyProvider | None = None, price_cap_per_bucket: int = 500, request_delay_sec: float = 2.0, concurrency: int = 4, resume_run_id: int | None = None, region_code: int = DEFAULT_REGION_CODE, ) -> YandexFullLoadCounters: """Exhaustive региональный сбор Yandex ЕКБ вторички (БЕЗ anchor'ов). Один проход с партиционированием по КОМНАТНОСТИ × ЦЕНЕ (адаптивное бинарное деление) охватывает весь ЕКБ. Инкрементальный save: on_bucket коммитит каждый leaf-бакет в БД сразу после сбора. Инжекция (#2135 F2): config/matcher/enrichment/shutdown_requested приходят снаружи вместо прямых импортов app.* (см. scraper_kit.contracts). record_yandex_price_history зовётся через enrichment (EnrichmentJobs Protocol). Cooperative cancel: runs.is_cancelled проверяется per-bucket. """ counters = YandexFullLoadCounters() skip_set: set[str] = set() if resume_run_id is not None: _prev_row = db.execute( text("SELECT counters FROM scrape_runs WHERE id = CAST(:rid AS bigint)"), {"rid": resume_run_id}, ).fetchone() if _prev_row is not None and _prev_row.counters: _prev_counters: dict = ( _prev_row.counters if isinstance(_prev_row.counters, dict) else {} ) skip_set = set(_prev_counters.get("done_buckets", [])) logger.info( "yandex-full-load run_id=%d: resuming from run %d — %d buckets already done", run_id, resume_run_id, len(skip_set), ) done: set[str] = set(skip_set) def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg] """Инкрементальный save после каждого leaf-бакета.""" nonlocal done if runs.is_cancelled(db, run_id): logger.info("yandex-full-load run_id=%d: cancel detected in on_bucket", run_id) raise RuntimeError("cancelled") elif shutdown_requested(): logger.info( "yandex-full-load run_id=%d: SIGTERM-drain — stopping at bucket %s", run_id, bucket_key, ) raise RuntimeError("shutdown") if not lots: done.add(bucket_key) runs.update_heartbeat( db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)} ) return inserted, updated = save_listings( db, lots, matcher=matcher, region_code=region_code, run_id=run_id, skip_seen_today=config.scraper_skip_seen_today, ) # save_listings вызывает db.commit() внутри — данные в БД сразу counters.saved_inserted += inserted counters.saved_updated += updated counters.unique_fetched += len(lots) # Price-history из gate price.previous/trend (graceful — провал истории # НЕ должен ронять сейв листингов/бакета). try: counters.price_history_rows += enrichment.record_yandex_price_history(db, lots) except Exception as ph_exc: logger.warning( "yandex-full-load run_id=%d: price-history failed bucket=%s: %s", run_id, bucket_key, ph_exc, ) try: db.rollback() except Exception: pass done.add(bucket_key) runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) logger.info( "yandex-full-load run_id=%d: bucket=%s saved ins=%d upd=%d total_unique=%d", run_id, bucket_key, inserted, updated, counters.unique_fetched, ) def _on_progress(unique_count: int) -> None: """Heartbeat per room-bucket (на случай пустых бакетов без on_bucket вызовов).""" runs.update_heartbeat(db, run_id, counters.to_dict()) try: async with YandexRealtyScraper(config, proxy_provider=proxy_provider) as scraper: scraper.request_delay_sec = request_delay_sec await scraper.fetch_all_secondary( price_cap_per_bucket=price_cap_per_bucket, concurrency=concurrency, on_bucket=_on_bucket, on_progress=_on_progress, skip_buckets=skip_set if skip_set else None, ) logger.info( "yandex-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d", run_id, counters.unique_fetched, counters.saved_inserted, counters.saved_updated, ) runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) logger.info( "yandex-full-load run_id=%d done: unique=%d ins=%d upd=%d errors=%d", run_id, counters.unique_fetched, counters.saved_inserted, counters.saved_updated, counters.errors_count, ) return counters except RuntimeError as exc: # on_bucket кидает RuntimeError("cancelled") при cooperative cancel, # RuntimeError("shutdown") при кооперативном SIGTERM-drain (#1182 Phase 3a). if str(exc) == "cancelled": logger.info( "yandex-full-load run_id=%d: cancelled — partial results unique=%d ins=%d upd=%d", run_id, counters.unique_fetched, counters.saved_inserted, counters.saved_updated, ) runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) return counters if str(exc) == "shutdown": logger.info( "yandex-full-load run_id=%d: SIGTERM-drain — partial results " "unique=%d ins=%d upd=%d", run_id, counters.unique_fetched, counters.saved_inserted, counters.saved_updated, ) runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) return counters logger.exception("yandex-full-load run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) raise except Exception as exc: logger.exception("yandex-full-load run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) raise # ── Avito exhaustive full load (citywide, без anchor'ов) ────────────────────── @dataclass class AvitoFullLoadCounters: """Aggregate counters для run_avito_full_load. Зеркало YandexFullLoadCounters — citywide room×price bisection без anchor/detail. """ unique_fetched: int = 0 saved_inserted: int = 0 saved_updated: int = 0 errors_count: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} async def run_avito_full_load( db: Session, *, run_id: int, config: ScraperConfig, matcher: HouseMatcher, shutdown_requested: Callable[[], bool] = lambda: False, price_cap_per_bucket: int = 1400, request_delay_sec: float = 7.0, concurrency: int = 5, secondary_only: bool = True, resume_run_id: int | None = None, incremental_days: int | None = None, region_code: int = DEFAULT_REGION_CODE, ) -> AvitoFullLoadCounters: """Exhaustive ИЛИ инкрементальный региональный сбор Avito ЕКБ вторички (БЕЗ anchor'ов). Один проход с партиционированием по КОМНАТНОСТИ × ЦЕНЕ (адаптивное бинарное деление ценового диапазона через `pmin`/`pmax`) охватывает весь ЕКБ. incremental_days: если задан — shallow инкрементальный обход (since-cutoff). Инжекция (#2135 F2): config/matcher/shutdown_requested приходят снаружи вместо прямых импортов app.* (см. scraper_kit.contracts). Cooperative cancel: runs.is_cancelled проверяется per-bucket. AvitoBlockedError/AvitoRateLimitedError → mark_banned (status='banned'). """ counters = AvitoFullLoadCounters() # ── Инкрементальный режим: since-cutoff для shallow обхода ──────────────── since: date | None = None if incremental_days is not None: since = date.today() - timedelta(days=incremental_days) logger.info( "avito-full-load run_id=%d: INCREMENTAL mode — since=%s (last %d days)", run_id, since.isoformat(), incremental_days, ) # ── Checkpoint/resume: читаем done_buckets из прошлого run ─────────────── skip_set: set[str] = set() if resume_run_id is not None: _prev_row = db.execute( text("SELECT counters FROM scrape_runs WHERE id = CAST(:rid AS bigint)"), {"rid": resume_run_id}, ).fetchone() if _prev_row is not None and _prev_row.counters: _prev_counters: dict = ( _prev_row.counters if isinstance(_prev_row.counters, dict) else {} ) skip_set = set(_prev_counters.get("done_buckets", [])) logger.info( "avito-full-load run_id=%d: resuming from run %d — %d buckets already done", run_id, resume_run_id, len(skip_set), ) done: set[str] = set(skip_set) def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg] """Инкрементальный save после каждого leaf-бакета.""" nonlocal done if runs.is_cancelled(db, run_id): logger.info("avito-full-load run_id=%d: cancel detected in on_bucket", run_id) raise RuntimeError("cancelled") elif shutdown_requested(): logger.info( "avito-full-load run_id=%d: SIGTERM-drain — stopping at bucket %s", run_id, bucket_key, ) raise RuntimeError("shutdown") if not lots: done.add(bucket_key) runs.update_heartbeat( db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)} ) return inserted, updated = save_listings( db, lots, matcher=matcher, region_code=region_code, run_id=run_id, skip_seen_today=config.scraper_skip_seen_today, ) # save_listings вызывает db.commit() внутри — данные в БД сразу counters.saved_inserted += inserted counters.saved_updated += updated counters.unique_fetched += len(lots) done.add(bucket_key) runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) logger.info( "avito-full-load run_id=%d: bucket=%s saved ins=%d upd=%d total_unique=%d", run_id, bucket_key, inserted, updated, counters.unique_fetched, ) def _on_progress(unique_count: int) -> None: """Heartbeat per room-bucket (на случай пустых бакетов без on_bucket вызовов).""" runs.update_heartbeat(db, run_id, counters.to_dict()) try: async with AvitoScraper(config) as scraper: scraper.request_delay_sec = request_delay_sec await scraper.fetch_all_secondary( price_cap_per_bucket=price_cap_per_bucket, concurrency=concurrency, secondary_only=secondary_only, on_bucket=_on_bucket, on_progress=_on_progress, skip_buckets=skip_set if skip_set else None, since=since, ) logger.info( "avito-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d", run_id, counters.unique_fetched, counters.saved_inserted, counters.saved_updated, ) runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) logger.info( "avito-full-load run_id=%d done: unique=%d ins=%d upd=%d errors=%d", run_id, counters.unique_fetched, counters.saved_inserted, counters.saved_updated, counters.errors_count, ) return counters except (AvitoBlockedError, AvitoRateLimitedError) as exc: # IP-бан/rate-limit: partial results уже сохранены через on_bucket. logger.error("avito-full-load run_id=%d: blocked/rate-limited — %s", run_id, exc) counters.errors_count += 1 runs.mark_banned( db, run_id, f"avito full load aborted: {exc}", {**counters.to_dict(), "done_buckets": sorted(done)}, ) return counters except RuntimeError as exc: # on_bucket кидает RuntimeError("cancelled") при cooperative cancel, # RuntimeError("shutdown") при кооперативном SIGTERM-drain (#1182 Phase 3a). if str(exc) == "cancelled": logger.info( "avito-full-load run_id=%d: cancelled — partial results unique=%d ins=%d upd=%d", run_id, counters.unique_fetched, counters.saved_inserted, counters.saved_updated, ) runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) return counters if str(exc) == "shutdown": logger.info( "avito-full-load run_id=%d: SIGTERM-drain — partial results " "unique=%d ins=%d upd=%d", run_id, counters.unique_fetched, counters.saved_inserted, counters.saved_updated, ) runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) return counters logger.exception("avito-full-load run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) raise except Exception as exc: logger.exception("avito-full-load run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) raise # ── DomClick citywide sweep (без anchor'ов) ──────────────────────────────── @dataclass class DomClickCitySweepCounters: """Aggregate counters для DomClick citywide sweep run. DomClick BFF — citywide (city_id), не geo-bbox: anchor-loop отсутствует. Один проход fetch_city перебирает ROOM_BUCKETS × pages внутри scraper'а. """ lots_fetched: int = 0 lots_inserted: int = 0 lots_updated: int = 0 pages_fetched: int = 0 errors_count: int = 0 blocked: int = 0 # 1 если QRATOR-блок был во время sweep geo_filtered: int = 0 # число офферов отфильтрованных geo-guard def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} # Дефолтные параметры sweep'а (EKB city_id=4). DOMCLICK_DEFAULT_CITY_ID: int = 4 DOMCLICK_DEFAULT_ROOMS: list[int] = [0, 1, 2, 3, 4] # vestigial; scraper sweeps all buckets # Число BFF-бакетов (st/1/2/3/4/5+) — фиксировано; используется для watchdog. _DOMCLICK_NUM_BUCKETS: int = 6 # Оценка времени одного fetch'а (network + parse) для watchdog. _DOMCLICK_PER_FETCH_S: float = 12.0 # Буфер сверху расчётного бюджета (cold browser start, save-фаза, price-splits). _DOMCLICK_SWEEP_BUDGET_S: float = 300.0 async def run_domclick_city_sweep( db: Session, *, run_id: int, config: ScraperConfig, matcher: HouseMatcher, shutdown_requested: Callable[[], bool] = lambda: False, city_id: int = DOMCLICK_DEFAULT_CITY_ID, rooms: list[int] | None = None, pages: int = 100, request_delay_sec: float | None = None, region_code: int = DEFAULT_REGION_CODE, ) -> DomClickCitySweepCounters: """DomClick citywide sweep через BFF JSON API. Структурно зеркалит run_cian_city_sweep / run_yandex_city_sweep, но CITYWIDE: DomClick не поддерживает geo-radius (fetch_around → NotImplementedError), поэтому anchor-loop отсутствует. DomClickScraper.fetch_city перебирает все ROOM_BUCKETS (st/1/2/3/4/5+) через BrowserFetcher → shared mobile proxy. Инжекция (#2135 F2): config/matcher/shutdown_requested приходят снаружи вместо прямых импортов app.* (см. scraper_kit.contracts). ЧЕСТНЫЙ СТАТУС (#1968): если scraper сообщил QRATOR-блок И lots == 0 → mark_failed. Иначе mark_done. Возвращает DomClickCitySweepCounters. """ _resolved_delay = request_delay_sec if request_delay_sec is not None else 6.0 counters = DomClickCitySweepCounters() # Watchdog: 6 buckets × pages × per_fetch + budget. _num_fetches = _DOMCLICK_NUM_BUCKETS * max(1, pages) _sweep_timeout = max( ANCHOR_TIMEOUT_SEC, int(_num_fetches * (_resolved_delay + _DOMCLICK_PER_FETCH_S) + _DOMCLICK_SWEEP_BUDGET_S), ) # Мутируемый контейнер для захвата scraper-ссылки из замыкания. _scraper_ref: list[DomClickScraper] = [] try: # ── Cooperative cancel перед SERP-фазой ────────────────────────────── if runs.is_cancelled(db, run_id): logger.info("domclick-sweep run_id=%d: cancelled before SERP", run_id) runs.update_heartbeat(db, run_id, counters.to_dict()) return counters elif shutdown_requested(): # SIGTERM-drain до SERP-фазы: финализируем run как mark_done(partial), # минуя honest-status (это чистый drain, не QRATOR-блок). logger.info("domclick-sweep run_id=%d: SIGTERM-drain — stopping before SERP", run_id) runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_done(db, run_id, counters.to_dict()) return counters logger.info( "domclick-sweep run_id=%d: BFF citywide sweep city_id=%d " "buckets=%d pages_cap=%d (watchdog %ds)", run_id, city_id, len(ROOM_BUCKETS), pages, _sweep_timeout, ) lots: list[ScrapedLot] = [] async def _domclick_phase() -> None: """Единственная citywide-фаза: fetch_city + save.""" nonlocal lots async with DomClickScraper(config) as _scraper: _scraper_ref.append(_scraper) if request_delay_sec is not None: _scraper.request_delay_sec = _resolved_delay lots = await _scraper.fetch_city(city_id=city_id, rooms=rooms, pages=pages) counters.lots_fetched += len(lots) if lots: inserted, updated = save_listings( db, lots, matcher=matcher, region_code=region_code, run_id=run_id ) counters.lots_inserted += inserted counters.lots_updated += updated try: await asyncio.wait_for(_domclick_phase(), timeout=_sweep_timeout) except TimeoutError: logger.warning( "domclick-sweep run_id=%d: SERP phase timed out after %ds — partial results", run_id, _sweep_timeout, ) counters.errors_count += 1 except Exception: logger.exception("domclick-sweep run_id=%d: SERP phase failed", run_id) counters.errors_count += 1 # Перенести счётчики scraper'а в counters (если scraper был создан). if _scraper_ref: _s = _scraper_ref[0] counters.blocked = 1 if _s.blocked else 0 counters.geo_filtered = _s.geo_filtered # Не-block fetch-ошибки скрейпера учитываем в errors_count. counters.errors_count += _s.fetch_errors # pages_fetched: worst-case число страниц (buckets × pages cap). counters.pages_fetched = _num_fetches runs.update_heartbeat(db, run_id, counters.to_dict()) # ── ЧЕСТНЫЙ СТАТУС (#1968) ──────────────────────────────────────────── if counters.lots_fetched == 0 and (counters.blocked or counters.errors_count > 0): logger.error( "domclick-sweep run_id=%d: 0 listings with blocked=%d errors=%d " "— marking failed", run_id, counters.blocked, counters.errors_count, ) runs.mark_failed( db, run_id, "QRATOR block or fetch errors — 0 listings", counters.to_dict(), ) else: runs.mark_done(db, run_id, counters.to_dict()) logger.info( "domclick-sweep run_id=%d done: lots=%d (ins=%d/upd=%d) " "pages=%d errors=%d blocked=%d geo_filtered=%d", run_id, counters.lots_fetched, counters.lots_inserted, counters.lots_updated, counters.pages_fetched, counters.errors_count, counters.blocked, counters.geo_filtered, ) return counters except Exception as exc: logger.exception("domclick-sweep run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) raise