Выбор оператора мобильного прокси опирался на две ненадёжные опоры. Первая: `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
361 lines
13 KiB
Python
361 lines
13 KiB
Python
"""Продуктовые адаптеры над `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",
|
||
]
|