gendesign/tradein-mvp/packages/scraper-kit/src/scraper_kit/contracts.py
bot-backend 964867a943
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI / changes (pull_request) Successful in 8s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 2m43s
feat(tradein/proxy): здоровье прокси по паре «узел × источник» (#2600 п.2)
Бан площадкой был глобальным: п.1 на распознанный бан выключал узел целиком
(enabled=false, disabled_reason='banned:<source>'). Реальность другая — Авито
банит IP, а Яндекс через тот же IP ходит чисто, поэтому один забаненный источник
выкидывал живой узел из пула для всех и худил пул быстрее, чем его пополняют
(#2638). Плюс такое состояние не самолечилось: ipify площадку не эмулирует, бан
не видит, а non-NULL disabled_reason блокирует авто-воскрешение (#2610) — нужен
был ручной PATCH.

Теперь бан — свойство ПАРЫ (proxy_id, source) в scrape_proxy_source_bans:
acquire(source) не выдаёт узел только этому источнику, для остальных узел
первосортный; снимается сам по времени. Срок эскалирует 6ч → 12 → 24 → 48 → 72
(потолок) на повторных банах той же пары; ban_count сбрасывается purge'ем
истёкших строк через 7 суток — поэтому purge намеренно отложенный, а не по
banned_until < now(). Защита последнего узла сохранена, но считается по
источнику: если после бана у acquire(source) не останется кандидатов — бан не
пишется, WARNING зовёт пополнять пул.

Миграция 210 конвертирует прод-остатки п.1 (enabled=false + disabled_reason
LIKE 'banned:%') в 6-часовые per-source баны и возвращает узлы в строй — иначе
они висели бы выключенными вечно.

Оператору активные баны видны в GET/PATCH /admin/proxies (source_bans) — без
этого «узел включён, но не выдаётся» необъяснимо.

Refs #2600
2026-08-05 17:14:12 +05:00

352 lines
18 KiB
Python
Raw 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.

"""Protocol-контракты (границы) между scraper_kit и продуктовым backend'ом.
Скрапперы (после переноса в scraper_kit, шаги C-G #2131) не должны напрямую
импортировать `app.*`. Вместо этого они принимают инжектируемые зависимости,
описанные здесь через `typing.Protocol` (structural typing): продуктовая сторона
предоставляет конкретные реализации (см. `app.services.scraper_adapters`),
которые структурно удовлетворяют этим Protocol'ам.
ВАЖНО: этот модуль НЕ импортирует `app.*` и не тянет тяжёлых зависимостей в
рантайме. Тип `Session` (SQLAlchemy) используется только в аннотациях под
`TYPE_CHECKING` — благодаря `from __future__ import annotations` аннотации не
вычисляются при импорте, поэтому пакет остаётся standalone-импортируемым.
Сигнатуры Protocol'ов зафиксированы по факту из backend-кода:
- HouseMatcher ← app.services.matching.houses.match_or_create_house
+ app.services.matching.listings.upsert_listing_source
- ScraperConfig ← поля app.core.config.Settings, читаемые скрапперами
- SessionFactory ← app.core.db.SessionLocal (sessionmaker → Session)
- EnrichmentJobs ← app.services.house_imv_backfill / house_dedup_merge /
yandex_address_backfill
"""
from __future__ import annotations
from collections.abc import Callable
from dataclasses import dataclass
from typing import TYPE_CHECKING, Any, Protocol, runtime_checkable
if TYPE_CHECKING:
from sqlalchemy.orm import Session
@dataclass(frozen=True)
class ProxyLease:
"""Арендованный из пула прокси (лёгкий транспорт между kit и backend).
Структурно совпадает с `app.services.proxy_pool.ProxyLease` (id/url/kind/rotate_url),
но живёт в контрактах, чтобы kit не импортировал `app.*`. Провайдеры читают только
`url` (готов для curl_cffi `proxies=`); `id`/`rotate_url` нужны backend-адаптеру для
release / mark_health / (будущей) ротации.
"""
id: int
url: str
kind: str
rotate_url: str | None = None
@runtime_checkable
class HouseMatcher(Protocol):
"""Кросс-источниковый матчинг домов/листингов (продуктовая логика matching-сервиса).
Реализуется `app.services.scraper_adapters.RealMatcherAdapter`, который проксирует
в `app.services.matching`. Скрапперы вызывают методы вместо прямого импорта.
"""
def match_or_create_house(
self,
db: Session,
ext_source: str,
ext_id: str,
address: str | None = ...,
lat: float | None = ...,
lon: float | None = ...,
*,
year_built: int | None = ...,
building_cadastral_number: str | None = ...,
cadastral_number: str | None = ...,
source_url: str | None = ...,
) -> tuple[int | None, float, str]:
"""Найти или создать канонический дом.
Returns:
(house_id, confidence ∈ [0.0, 1.0], method), где method ∈ {
'cadastr_exact', 'source_exact', 'fingerprint', 'geo_proximity', 'new',
'no_house_number'
}.
house_id=None только при method='no_house_number' — matcher отказался
создавать дом из безномерного адреса без кадастра (P1). Вызывающий обязан
обработать None.
"""
...
def upsert_listing_source(
self,
db: Session,
*,
listing_id: int,
ext_source: str,
ext_id: str,
method: str,
confidence: float,
price_rub: float | None = ...,
area_m2: float | None = ...,
floor: int | None = ...,
rooms_count: int | None = ...,
source_url: str | None = ...,
source_data: dict[str, Any] | None = ...,
) -> None:
"""Зарегистрировать/обновить строку listing_sources для (ext_source, ext_id)."""
...
@runtime_checkable
class ScraperConfig(Protocol):
"""Конфигурация, читаемая скрапперами из `app.core.config.settings`.
Поля собраны grep'ом `settings.<name>` по `app/services/scrapers/`. Типы —
как в `Settings` (proxy-свойства возвращают `str | None`). Реализуется
`app.services.scraper_adapters.RealScraperConfig` (проксирует в singleton
`settings`); pydantic-`Settings` также структурно удовлетворяет этому Protocol.
"""
# Режим загрузки: "curl_cffi" | "browser".
scraper_fetch_mode: str
# HTTP-эндпоинт headless-браузера (tradein-browser).
browser_http_endpoint: str
# Backconnect-прокси общий (Avito/Cian/Yandex, #2616 шаг 2: единственный источник).
scraper_proxy_url: str | None
# Лимит IP-ротаций на прогон (changeip-ссылка снята #2616 шаг 2 — сама ротация
# no-op, поле гейтит budget accounting в ban-rotation state machine).
avito_proxy_max_rotations: int
# Ограничивать Avito SERP только ЕКБ.
avito_serp_ekb_only: bool
# Cian-прокси (#2616 шаг 2: = scraper_proxy_url, per-provider override снят).
cian_proxy_url: str | None
# Границы валидности cian-оценки, руб.
cian_valuation_min_rub: float
cian_valuation_max_rub: float
# Секрет pgp_sym_encrypt для cian_session_cookies (read в session.py).
cookie_encryption_key: str
# Sentry/GlitchTip DSN (опционально).
glitchtip_dsn: str | None
# ── Оркестрация sweep-pipeline (#2135) ───────────────────────────────────
# Поля, читаемые слоем оркестрации (scraper_kit.orchestration.pipeline).
# Собраны grep'ом `settings.<name>` по app/services/scrape_pipeline.py.
# Проксируются той же продуктовой реализацией (RealScraperConfig).
#
# Partial-ban intake: SERP уже собрал лоты, а detail/houses заблокировало —
# run помечается 'done' (не 'banned'), intake сохранён (#1950).
avito_serp_ok_not_banned: bool
# changeip settle-пауза (секунды) после ротации IP.
avito_proxy_rotate_settle_s: float
# Ретраи changeip-GET: N попыток по proxy_rotate_attempt_timeout_s каждая (#1950).
# changeip-ссылка снята (#2616 шаг 2) — эти два поля сейчас не читаются
# ban-rotation кодом (см. avito_proxy_max_rotations выше), оставлены как
# budget-верхняя-граница для app.tasks.avito_detail_backfill wait_for.
proxy_rotate_attempts: int
proxy_rotate_attempt_timeout_s: float
# Cian лимит ротаций (rotate_url снят #2616 шаг 2, см. avito_proxy_max_rotations).
cian_proxy_max_rotations: int
# Yandex лимит ротаций (аналогично).
yandex_proxy_max_rotations: int
# Пропускать листинги, уже обновлённые сегодня по МСК (full-load дедуп).
scraper_skip_seen_today: bool
# Per-fetch watchdog-таймаут Cian full-load (секунды).
cian_full_load_per_fetch_timeout_s: float
# ── Proxy-pool (#2163) ────────────────────────────────────────────────────
# Feature-флаг: брать прокси из пула (scrape_proxies) на curl_cffi-путях вместо
# env-прокси (config.*_proxy_url). Ship-dark: дефолт False → curl-пути ходят как
# сейчас. True + пустой пул → fallback на env (сбор не ломается). См. ProxyProvider.
use_proxy_pool_curl: bool
# Feature-флаг (#2164 P4): брать прокси из пула на browser-путях (BrowserFetcher →
# POST /fetch{proxy}). Scraper-сторона делает acquire(source) и передаёт lease.url в
# теле запроса; tradein-browser релончит camoufox с этим прокси (relaunch ТОЛЬКО при
# реальной смене). Ship-dark: дефолт False → браузер берёт прокси из env
# (BROWSER_PROXY_*), прод не меняется. True + пустой пул → fallback на env.
use_proxy_pool_browser: bool
# ── #2616 шаг 1: признак окружения для отказа вместо мёртвого env-fallback ──────
# "production" в прод-контейнерах (ENV: ENVIRONMENT), иначе dev/test/local. Читают
# curl_proxy_url (_proxy.py) и BrowserFetcher._pool_proxy — пул пуст/acquire упал +
# окружение НЕ "production" → легитимный dev/no-op fallback (без изменений).
# environment == "production" → NoProxyAvailableError вместо мёртвого env-прокси.
environment: str
@runtime_checkable
class ProxyProvider(Protocol):
"""Источник прокси из пула для curl-путей (#2163), за флагом use_proxy_pool_curl.
Инкапсулирует lease/health-механику `app.services.proxy_pool` (которая требует
`Session`), пряча БД от kit'а. Реализуется
`app.services.scraper_adapters.RealProxyProvider` (открывает короткую сессию на
операцию). Провайдеры вызывают `acquire` → используют `lease.url` → на выходе
`mark_health` + `release` (в finally — lease не должен течь).
Контракт ship-dark/fallback: `acquire` возвращает None если свободных здоровых
прокси нет — caller обязан откатиться на env-прокси, а НЕ падать.
"""
def acquire(self, provider: str) -> ProxyLease | None:
"""Взять прокси под провайдера (avito/cian/yandex/generic/any).
Returns ProxyLease или None если пул пуст/выключен → caller делает env-fallback.
"""
...
def release(self, lease: ProxyLease) -> None:
"""Освободить lease (идемпотентно). Вызывать в finally — иначе прокси течёт."""
...
def mark_health(
self,
lease: ProxyLease,
ok: bool,
*,
exit_ip: str | None = ...,
latency_ms: int | None = ...,
) -> None:
"""Записать исход использования прокси (ok=True успех, False бан/ошибка)."""
...
def mark_banned(self, lease: ProxyLease, *, source: str) -> None:
"""Пометить lease забаненным площадкой `source` (#2600 п.1, п.2).
Отличается от `mark_health(ok=False)`: та инкрементит consecutive_fails и
авто-disable'ит только после DISABLE_THRESHOLD подряд неудач (транзиентный
сбой должен пережить пару неудач). Здесь сигнал УЖЕ надёжно распознан (валидная
HTML-заглушка/капча/QRATOR-маркер, не сетевая ошибка) — узел немедленно снимается
с выдачи ЭТОМУ источнику (строка в `scrape_proxy_source_bans`, срок эскалирует
на повторных банах), для остальных источников остаётся в строю: площадка банит
IP, а не ломает прокси. Исключение — последний узел, достижимый для `source`:
бан не записывается (см. `app.services.proxy_pool.mark_banned`, там же защита).
Вызывать из точки детекта бана, ПОКА lease ещё держится (до release/__aexit__) —
иначе id узла, который был использован, потерян. Best-effort — caller
(`BrowserFetcher.report_ban` / `curl_proxy_url`) не должен падать на ошибке пула.
"""
...
def touch(self, lease: ProxyLease) -> None:
"""Heartbeat: продлить lease (без изменения health-полей).
Для долгоживущих browser-сессий (`BrowserFetcher` держит ОДИН lease на весь
прогон, часы) — вызывать периодически (на каждый /fetch), чтобы серверный
`reap_stale_leases` не отобрал активно используемый прокси у прогона дольше
`STALE_LEASE_MINUTES`. Best-effort — caller не должен падать при ошибке.
"""
...
@runtime_checkable
class SessionFactory(Protocol):
"""Фабрика SQLAlchemy-сессии (аналог `app.core.db.SessionLocal`).
Вызов без аргументов возвращает новую `Session`. Реализуется
`app.services.scraper_adapters.RealSessionFactory`; сам `sessionmaker`
(SessionLocal) также структурно удовлетворяет.
"""
def __call__(self) -> Session:
"""Открыть новую сессию БД."""
...
@runtime_checkable
class EnrichmentJobs(Protocol):
"""Реестр enrichment-джоб, дёргаемых после сохранения скрапом.
Callback'и проксируют в соответствующие `app.services.*`. Возвращаемые
доменные результат-типы (`HouseIMVBackfillResult`, `YandexAddressBackfillResult`)
остаются `Any`, чтобы не тянуть `app.*` в контракты. Реализуется
`app.services.scraper_adapters.RealEnrichmentJobs`.
"""
async def house_imv_backfill(
self,
db: Session,
*,
batch_size: int = ...,
request_delay_sec: float = ...,
only_status: str = ...,
house_id: int | None = ...,
heartbeat: Callable[[], None] | None = ...,
) -> Any:
"""Avito IMV-оценка для домов в scope. Проксирует backfill_house_imv."""
...
def house_dedup_merge(
self,
db: Session,
*,
run_id: int,
params: dict[str, Any],
) -> dict[str, int]:
"""Слить дубли домов. Проксирует run_house_dedup_merge. Returns счётчики."""
...
async def yandex_address_backfill(
self,
db: Session,
*,
limit: int = ...,
request_delay_sec: float = ...,
) -> Any:
"""Догрузить точные адреса Yandex. Проксирует backfill_yandex_addresses."""
...
async def process_houses_imv_batch(
self,
db: Session,
house_ids: set[int],
*,
request_delay_sec: float = ...,
heartbeat: Callable[[], None] | None = ...,
) -> Any:
"""Avito IMV-оценка явного набора house_id (финальная IMV-фаза city-sweep).
Проксирует process_houses_imv_batch. Возвращаемый результат имеет поля
.checked/.saved/.errors (тип остаётся Any — не тянем app.* в контракты).
"""
...
# ── Yandex sweep-хелперы (#2135 F2) ──────────────────────────────────────
# Продуктовая логика, дёргаемая yandex city/full sweep-оркестраторами напрямую
# (не через отдельные backfill-джобы). Инжектируются вместо прямых импортов
# app.services.yandex_price_history / yandex_address_backfill.
def record_yandex_price_history(self, db: Session, lots: list[Any]) -> int:
"""Записать price-history из gate price.previous/trend. Проксирует
record_yandex_price_history. Returns число вставленных строк истории.
`lots` — список kit-`ScrapedLot` (тип оставлен `Any`, чтобы не тянуть
доменный тип в контракт). Читается duck-typing'ом (source_id/price/…).
"""
...
def extract_address_from_title(self, html: str) -> str | None:
"""Извлечь полный адрес из <title> detail-страницы Yandex.
Проксирует `_extract_address_from_title` (yandex_address_backfill).
Returns адрес с номером дома или None если не распарсилось.
"""
...
def address_has_house_number(self, address: str) -> bool:
"""True если адрес содержит номер дома (`,\\s*\\d+`).
Проксирует regex `_RE_HAS_HOUSE_NUMBER` из yandex_address_backfill —
gate для решения, обновлять ли address на обогащённый.
"""
...
__all__ = [
"EnrichmentJobs",
"HouseMatcher",
"ProxyLease",
"ProxyProvider",
"ScraperConfig",
"SessionFactory",
]