gendesign/tradein-mvp/backend/app/services/scraper_adapters.py
bot-backend b89788ee99 feat(tradein/proxy): прогон знает свой узел, а снятый бан перестаёт стирать историю (#3404)
Выбор оператора мобильного прокси опирался на две ненадёжные опоры.

Первая: `scrape_runs` не знала, через какой узел шёл прогон — колонка `proxy_id`
была только у банов и ротаций. «Какой узел собрал 5 карточек из 21» не выяснялось
ни одним запросом.

Вторая: `clear_source_bans` делала DELETE, а зовётся она после КАЖДОЙ успешной
ротации exit-IP. У #540723 (МегаФон) 23 успешные ротации и ноль строк банов,
у #540722 (Tele2) ротаций почти не было и 7 банов. «7 против 0» читалось как
«Tele2 хуже», хотя в той же мере это «у МегаФона историю стёрли 23 раза».

Теперь:
- `scrape_runs.proxy_id` — последний выданный прогону узел; полная цепочка
  (если узел менялся mid-run) копится в `counters.proxy_ids`. Пишет
  `proxy_pool.attribute_run_proxy` из единственной точки — сразу после выдачи
  лиза в `acquire()`, поэтому curl-путь, браузерный sticky lease и ре-acquire
  при ротации покрыты одинаково. `run_id` доходит до адаптера через ContextVar
  (`scraper_kit.orchestration.run_context`): протокол `ProxyProvider.acquire`
  его не несёт, а `RealProxyProvider` живёт одним объектом на весь планировщик.
  Best-effort: `lock_timeout` 2с и проглоченное исключение — диагностика не
  вправе ронять выдачу прокси или ждать на блокировке строки прогона.
- `clear_source_bans` гасит строку (`banned_until = now()`, `ban_count = 0`,
  `cleared_at`/`cleared_reason`) вместо удаления. Эскалация сохраняется 1:1:
  формула в `mark_banned` берёт ПРЕДЫДУЩИЙ `ban_count` показателем степени, при
  нуле это ровно `SOURCE_BAN_BASE_HOURS` — как после DELETE. Строка доживает до
  штатного purge по `SOURCE_BAN_PURGE_DAYS`.

Для всех читателей `scrape_proxy_source_bans` погашенная строка неотличима от
отсутствующей: acquire, оба guard-подзапроса `mark_banned`, `proxy_egress`
(ранжирование по `ban_count` даёт 0, как у узла без истории), admin `_active_ban` —
все гейтятся по `banned_until > now()`.

Ничего не бэкфиллится: связать прошедшие прогоны с узлами нечем (`leased_by`
исторически = NON_RUN_LEASE_MARKER), врать восстановленным значением нельзя.

Миграция 287. Тесты: 9 новых на обе части (главный — эскалация после гашения даёт
базовые 6ч, а не удвоенные) + 14 существующих переведены с DELETE-семантики на
гашение, включая проверку, что секрет ротации не утекает в новое `cleared_reason`.
Полный прогон бэкенда: 5600 passed, 37 skipped.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_011WHFxVPWoBnSZihkdH1Uou
2026-09-06 13:05:20 +03:00

361 lines
13 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

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

"""Продуктовые адаптеры над `app.services.*`, удовлетворяющие scraper_kit.contracts.
Тонкие врапперы, инжектируемые в скрапперы (после переноса в scraper_kit, шаги C-G
#2131) вместо прямых импортов `app.*`. Каждый класс структурно удовлетворяет
соответствующему Protocol'у из `scraper_kit.contracts`:
RealMatcherAdapter → HouseMatcher
RealScraperConfig → ScraperConfig
RealSessionFactory → SessionFactory
RealEnrichmentJobs → EnrichmentJobs
Wired в прод с #2192: `app.scheduler_main._run_kit_scheduler` конструирует
`SchedulerContext` из этих адаптеров (`RealScraperConfig`/`RealMatcherAdapter`/
`RealEnrichmentJobs`/`RealSessionFactory`/`RealProxyProvider`) и передаёт в
`scraper_kit.orchestration.scheduler.scheduler_loop` — kit-native sweep-registry
крутится через эту инжекцию, а не через прямые импорты `app.*`. Держим сигнатуры
в точном соответствии с исходными функциями, чтобы Protocol'ы `scraper_kit.contracts`
не расходились с боевой реализацией.
"""
from __future__ import annotations
from collections.abc import Callable
from typing import TYPE_CHECKING, Any
from app.core.config import settings as _settings
from app.core.db import SessionLocal as _SessionLocal
from app.services import proxy_pool as _proxy_pool
from app.services.house_dedup_merge import run_house_dedup_merge as _run_house_dedup_merge
from app.services.house_imv_backfill import backfill_house_imv as _backfill_house_imv
from app.services.house_imv_backfill import (
process_houses_imv_batch as _process_houses_imv_batch,
)
from app.services.matching import match_or_create_house as _match_or_create_house
from app.services.matching import upsert_listing_source as _upsert_listing_source
from app.services.yandex_address_backfill import (
_RE_HAS_HOUSE_NUMBER,
_extract_address_from_title,
)
from app.services.yandex_address_backfill import (
backfill_yandex_addresses as _backfill_yandex_addresses,
)
from app.services.yandex_price_history import (
record_yandex_price_history as _record_yandex_price_history,
)
if TYPE_CHECKING:
from scraper_kit.contracts import ProxyLease
from sqlalchemy.orm import Session
from app.services.house_imv_backfill import HouseIMVBackfillResult
from app.services.yandex_address_backfill import YandexAddressBackfillResult
class RealMatcherAdapter:
"""HouseMatcher-адаптер над `app.services.matching`."""
def match_or_create_house(
self,
db: Session,
ext_source: str,
ext_id: str,
address: str | None = None,
lat: float | None = None,
lon: float | None = None,
*,
year_built: int | None = None,
building_cadastral_number: str | None = None,
source_url: str | None = None,
city: str | None = None,
region_code: int | None = None,
) -> tuple[int | None, float, str]:
# house_id is None when the matcher refuses a numberless address without a
# cadastral number (method 'no_house_number', P1). Callers must tolerate None.
return _match_or_create_house(
db,
ext_source,
ext_id,
address,
lat,
lon,
year_built=year_built,
building_cadastral_number=building_cadastral_number,
source_url=source_url,
city=city,
region_code=region_code,
)
def upsert_listing_source(
self,
db: Session,
*,
listing_id: int,
ext_source: str,
ext_id: str,
method: str,
confidence: float,
price_rub: float | None = None,
area_m2: float | None = None,
floor: int | None = None,
rooms_count: int | None = None,
source_url: str | None = None,
source_data: dict[str, Any] | None = None,
) -> None:
_upsert_listing_source(
db,
listing_id=listing_id,
ext_source=ext_source,
ext_id=ext_id,
method=method,
confidence=confidence,
price_rub=price_rub,
area_m2=area_m2,
floor=floor,
rooms_count=rooms_count,
source_url=source_url,
source_data=source_data,
)
class RealScraperConfig:
"""ScraperConfig-адаптер: read-only проекция singleton'а `settings`.
Свойства проксируют в `app.core.config.settings`, включая computed-properties
(`scraper_proxy_url`, `cian_proxy_url`), чтобы скрапперы читали ровно те же
значения, что и сейчас через прямой `settings.<name>`.
"""
@property
def scraper_fetch_mode(self) -> str:
return _settings.scraper_fetch_mode
@property
def browser_http_endpoint(self) -> str:
return _settings.browser_http_endpoint
@property
def scraper_proxy_url(self) -> str | None:
return _settings.scraper_proxy_url
@property
def avito_proxy_max_rotations(self) -> int:
return _settings.avito_proxy_max_rotations
@property
def avito_serp_ekb_only(self) -> bool:
return _settings.avito_serp_ekb_only
@property
def cian_proxy_url(self) -> str | None:
return _settings.cian_proxy_url
@property
def cian_valuation_min_rub(self) -> float:
return _settings.cian_valuation_min_rub
@property
def cian_valuation_max_rub(self) -> float:
return _settings.cian_valuation_max_rub
@property
def cookie_encryption_key(self) -> str:
return _settings.cookie_encryption_key
@property
def glitchtip_dsn(self) -> str | None:
return _settings.glitchtip_dsn
# ── Оркестрация sweep-pipeline (#2135) ───────────────────────────────
@property
def avito_serp_ok_not_banned(self) -> bool:
return _settings.avito_serp_ok_not_banned
@property
def avito_proxy_rotate_settle_s(self) -> float:
return _settings.avito_proxy_rotate_settle_s
@property
def cian_proxy_max_rotations(self) -> int:
return _settings.cian_proxy_max_rotations
@property
def yandex_proxy_max_rotations(self) -> int:
return _settings.yandex_proxy_max_rotations
@property
def scraper_skip_seen_today(self) -> bool:
return _settings.scraper_skip_seen_today
@property
def cian_full_load_per_fetch_timeout_s(self) -> float:
return _settings.cian_full_load_per_fetch_timeout_s
# ── Proxy-pool (#2163) ────────────────────────────────────────────────
@property
def use_proxy_pool_curl(self) -> bool:
return _settings.use_proxy_pool_curl
# ── Proxy-pool browser (#2164 P4) ─────────────────────────────────────
@property
def use_proxy_pool_browser(self) -> bool:
return _settings.use_proxy_pool_browser
# ── #2616 шаг 1: признак окружения для отказа вместо мёртвого env-fallback ──
@property
def environment(self) -> str:
return _settings.environment
class RealProxyProvider:
"""ProxyProvider-адаптер над `app.services.proxy_pool` (#2163).
Прячет `Session` от kit'а: открывает короткую сессию на каждую операцию
(acquire/release/mark_health), проксирует в модульные функции proxy_pool и
коммитит внутри них. Возвращает kit-`ProxyLease` (структурно совпадает с
`proxy_pool.ProxyLease`), чтобы провайдеры не импортировали `app.*`.
Не-run вызовы (curl-путь эстиматора вне scrape_run) лизят под
NON_RUN_LEASE_MARKER — reap_stale_leases освободит, если release не дошёл.
"""
def acquire(self, provider: str) -> ProxyLease | None:
from scraper_kit.contracts import ProxyLease as _KitProxyLease
from scraper_kit.orchestration.run_context import current_run_id
# #3404: протокол ProxyProvider.acquire(provider) не несёт run_id (этот
# адаптер — один объект на весь scheduler_main.py), поэтому берём его из
# ContextVar, который выставляет runs.create_run. None — вызов вне прогона
# (health-check, эстиматор) — proxy_pool.acquire в этом случае лизит под
# NON_RUN_LEASE_MARKER, как и раньше, атрибуцию в scrape_runs не пишет.
run_id = current_run_id.get()
db = _SessionLocal()
try:
lease = _proxy_pool.acquire(db, provider, run_id=run_id)
finally:
db.close()
if lease is None:
return None
return _KitProxyLease(
id=lease.id, url=lease.url, kind=lease.kind, rotate_url=lease.rotate_url
)
def release(self, lease: ProxyLease) -> None:
db = _SessionLocal()
try:
_proxy_pool.release(db, lease.id)
finally:
db.close()
def mark_health(
self,
lease: ProxyLease,
ok: bool,
*,
exit_ip: str | None = None,
latency_ms: int | None = None,
) -> None:
db = _SessionLocal()
try:
_proxy_pool.mark_health(db, lease.id, ok, exit_ip=exit_ip, latency_ms=latency_ms)
finally:
db.close()
def touch(self, lease: ProxyLease) -> None:
db = _SessionLocal()
try:
_proxy_pool.touch(db, lease.id)
finally:
db.close()
def mark_banned(self, lease: ProxyLease, *, source: str) -> None:
db = _SessionLocal()
try:
_proxy_pool.mark_banned(db, lease.id, source=source)
finally:
db.close()
class RealSessionFactory:
"""SessionFactory-адаптер над `app.core.db.SessionLocal`."""
def __call__(self) -> Session:
return _SessionLocal()
class RealEnrichmentJobs:
"""EnrichmentJobs-адаптер над enrichment-сервисами `app.services.*`."""
async def house_imv_backfill(
self,
db: Session,
*,
batch_size: int = 50,
request_delay_sec: float = 5.0,
only_status: str = "pending",
house_id: int | None = None,
heartbeat: Callable[[], None] | None = None,
) -> HouseIMVBackfillResult:
return await _backfill_house_imv(
db,
batch_size=batch_size,
request_delay_sec=request_delay_sec,
only_status=only_status,
house_id=house_id,
heartbeat=heartbeat,
)
def house_dedup_merge(
self,
db: Session,
*,
run_id: int,
params: dict[str, Any],
) -> dict[str, int]:
return _run_house_dedup_merge(db, run_id=run_id, params=params)
async def yandex_address_backfill(
self,
db: Session,
*,
limit: int = 200,
request_delay_sec: float = 3.0,
) -> YandexAddressBackfillResult:
return await _backfill_yandex_addresses(
db,
limit=limit,
request_delay_sec=request_delay_sec,
)
async def process_houses_imv_batch(
self,
db: Session,
house_ids: set[int],
*,
request_delay_sec: float = 5.0,
heartbeat: Callable[[], None] | None = None,
) -> HouseIMVBackfillResult:
return await _process_houses_imv_batch(
db,
house_ids,
request_delay_sec=request_delay_sec,
heartbeat=heartbeat,
)
def record_yandex_price_history(self, db: Session, lots: list[Any]) -> int:
return _record_yandex_price_history(db, lots)
def extract_address_from_title(self, html: str) -> str | None:
return _extract_address_from_title(html)
def address_has_house_number(self, address: str) -> bool:
return bool(_RE_HAS_HOUSE_NUMBER.search(address))
__all__ = [
"RealEnrichmentJobs",
"RealMatcherAdapter",
"RealProxyProvider",
"RealScraperConfig",
"RealSessionFactory",
]