gendesign/tradein-mvp/backend/app/api/v1/admin.py
bot-backend 70f2daf241
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 9s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 3m54s
fix(tradein/cian): обогащение ЖК падало не на разметке, а на сожжённом узле (#2767)
Живая проба 2026-08-09 тем же трактом (сайдкар → camoufox → узел пула), одна и та
же страница zhk-pihtovyy-ekb-i.cian.ru:

    узел #1  asocks-residential-1 → 374 168 Б, cian_waf_block, initialState нет
    узел #9  asocks-mobile-1      → 542 021 Б,                 initialState ЕСТЬ
    узел #10 asocks-mobile-2      → 557 235 Б,                 initialState ЕСТЬ
    узел #11 asocks-mobile-3      → 557 163 Б,                 initialState ЕСТЬ

Разметка не менялась: MFE 'newbuilding-card-desktop-frontend'/'initialState'
разбирается ТЕКУЩИМ кодом на трёх узлах из четырёх. Страница в 374 КБ — не
карточка и не SPA-оболочка, а страница блокировки Циана: «Обнаружен
подозрительный трафик», код страницы cian_waf_block, собственная аналитика
помечает её pageType:"VPNBlock".

Почему это восемь суток читалось как смена разметки: fetch_newbuilding была
единственным cian-путём, который строил BrowserFetcher вручную, без
proxy_provider (serp.py всегда шёл через build_browser_fetcher). Значит сбор
всегда выходил через ОДИН env-узел сайдкара — SCRAPER_PROXY_URL, он же узел
пула #1, чей exit-IP Циан забанил. Ротация была невозможна, report_ban без
lease — no-op, 25 попыток за ночь уходили в тот же адрес.

Диагностика #2768 при этом отвечала «antibot_markers=none»: список маркеров не
знал подписи cian_waf_block, а страница блока не содержит ни «captcha», ни «вы
не робот». Молчание списка прочиталось как его вердикт «защиты нет» — и увело
разбор в гипотезу про разметку.

Правки:
- fetch_newbuilding принимает proxy_provider и строится через
  build_browser_fetcher (пул, как у всех остальных cian-путей);
- страница блока репортится в пул как бан узла ДО выхода из контекста, пока
  lease жив, — дальше acquire('cian') этот узел не выдаёт;
- cian_waf_block добавлен в _ANTIBOT_MARKERS; размер сам по себе гипотезы
  больше не разделяет (блок весит 374 КБ) — это записано в docstring;
- proxy_provider прокинут во все три вызывающих: обогащение, cian_history,
  admin debug-роут;
- из лога убрано «(captcha / parse miss?)» — догадка автора кода, читавшаяся
  дальше как факт.

Тест на НАСТОЯЩЕМ сохранённом ответе прода (фикстура cian_waf_block_zhk_page.html,
обрезаны инлайн-стили). На коде до правки: _describe_parse_miss даёт
«antibot_markers=none», _blocked_by нет, аргумента proxy_provider нет.

Refs #2767
2026-08-09 22:14:20 +05:00

3113 lines
130 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.

"""Admin endpoints — debug + ручной запуск парсеров.
Auth: Caddy basic_auth gate'ит /trade-in/api/v1/admin/* — application layer открыт.
"""
from __future__ import annotations
import asyncio
import json
import logging
import time
from typing import Annotated, Any, Literal
from urllib.parse import urlparse, urlunparse
from uuid import uuid4
import httpx
from fastapi import APIRouter, BackgroundTasks, Depends, HTTPException, Query
from pydantic import BaseModel, Field, field_validator
# scraper_kit-миграция (#2305, Group A): debug run-parser роуты этого файла переключены
# на scraper_kit.providers.* эквиваленты (парность доказана tests/scrapers/parity — см.
# vault Scraper_Kit_Legacy_Dependency_Audit_0703). YandexNewbuildingScraper.fetch_jk и
# cian_newbuilding.fetch_newbuilding были временно исключены — их kit-эквиваленты
# (providers/yandex/newbuilding.py, providers/cian/newbuilding.py) строили
# BrowserFetcher(source=...) без обязательного kwarg'а endpoint → TypeError на любом
# вызове. Исправлено в #2322 (endpoint пробрасывается через ScraperConfig,
# config=RealScraperConfig() передаётся явно) — оба debug-роута этого файла
# (#2397 Part D4) теперь тоже на kit.
#
# scraper_kit-миграция (#2397 slice A, эпик #2277): 5 debug city-sweep/full-load
# роутов (avito-city-sweep, cian-city-sweep, cian-full-load, yandex-city-sweep,
# yandex-full-load) переключены с app.services.scrape_pipeline (legacy, удалён в
# #2397 Part E1) на scraper_kit.orchestration.pipeline — тот же DI-паттерн, что и
# app.scheduler_main._run_kit_scheduler / scraper_kit.orchestration.scheduler
# ._job_*: config/matcher/enrichment/proxy_provider инжектируются явно вместо
# module-level импортов app.*. shutdown_requested НЕ прокидывается (нет SIGTERM-drain
# семантики у admin BackgroundTasks — эквивалент дефолту lambda: False, поведение не
# меняется). Боевой scheduler-путь (production sweep-cron) уже на kit (#2397 Part C
# убрал legacy app.services.scheduler.scheduler_loop fallback) — тот же orchestrator,
# что и debug-роуты этого файла.
from scraper_kit.base import save_listings
from scraper_kit.browser_fetcher import BrowserFetcher
from scraper_kit.orchestration.pipeline import (
DEFAULT_REGION_CODE,
run_avito_city_sweep,
run_cian_city_sweep,
run_cian_full_load,
run_yandex_city_sweep,
run_yandex_full_load,
)
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.imv import (
IMVAddressNotFoundError,
evaluate_via_imv,
save_imv_evaluation,
save_imv_placement_history,
)
from scraper_kit.providers.avito.serp import AvitoScraper
from scraper_kit.providers.cian.serp import CianScraper
from scraper_kit.providers.yandex.detail import YandexDetailScraper
from scraper_kit.providers.yandex.newbuilding import YandexNewbuildingScraper
from scraper_kit.providers.yandex.serp import YandexRealtyScraper
from scraper_kit.providers.yandex.valuation import YandexValuationScraper
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.core.config import settings
from app.core.db import SessionLocal, get_db
from app.schemas.trade_in import ScheduleConfig, ScheduleConfigUpdate
from app.services import cian_session as cian_session_svc
from app.services import domclick_session as domclick_session_svc
from app.services import proxy_rotation as proxy_rotation_svc
from app.services import scrape_runs as runs_mod
from app.services.estimator import LISTINGS_FRESH_DAYS
from app.services.geocoder import geocode, known_city_hint
from app.services.proxy_pool import clear_source_bans
from app.services.scheduler import has_running_run
from app.services.scraper_adapters import (
RealEnrichmentJobs,
RealMatcherAdapter,
RealProxyProvider,
RealScraperConfig,
)
from app.services.scraper_settings import get_scraper_delay, invalidate_cache
from app.tasks.avito_detail_backfill import run_avito_detail_backfill
from app.tasks.geocode_missing import GeocodeBackfillResult, geocode_missing_listings
logger = logging.getLogger(__name__)
router = APIRouter()
def _kit_proxy_provider() -> RealProxyProvider | None:
"""RealProxyProvider только когда proxy-pool включён (#2163/#2164) — иначе None.
Зеркалит гейт `app.scheduler_main._run_kit_scheduler`: оба флага дефолтно False
(ship-dark), поэтому debug-роуты этого файла ведут себя как раньше (env-прокси),
пока пул явно не включён.
"""
if settings.use_proxy_pool_curl or settings.use_proxy_pool_browser:
return RealProxyProvider()
return None
def _assert_allowed_url(url: str) -> None:
"""SSRF guard: отклоняем абсолютные URL с хостом не из allowlist (#756).
Относительные пути (urlparse().netloc пустой) пропускаем — хост
будет подставлен фиксированным AVITO_BASE / BASE_URL в самом scraper'е.
"""
parsed = urlparse(url)
if parsed.netloc and parsed.netloc not in settings.scrape_allowed_hosts:
raise HTTPException(status_code=400, detail="host not allowed")
def _assert_domclick_host(url: str) -> None:
"""SSRF guard для DomClick card-URL — вариант _assert_allowed_url для поддоменов.
settings.scrape_allowed_hosts содержит только точные "domclick.ru"/"www.domclick.ru",
а реальные карточки живут на региональных поддоменах (ekaterinburg.domclick.ru,
spb.domclick.ru, ...) — точное сравнение из _assert_allowed_url их бы отклонило.
Разрешаем domclick.ru и любой *.domclick.ru поддомен; относительные пути
(netloc пуст) пропускаем как и _assert_allowed_url.
"""
parsed = urlparse(url)
host = parsed.netloc.lower()
if host and host != "domclick.ru" and not host.endswith(".domclick.ru"):
raise HTTPException(status_code=400, detail="host not allowed")
def _safe_http_error(status_code: int, public_detail: str, log_context: str) -> HTTPException:
"""Defense-in-depth: не отдаём текст исключения в HTTP-ответ.
Текст psycopg/httpx-исключений может содержать DSN-фрагменты, внутренние
хосты, пути. Логируем полный exception с коротким error-id через
``logger.exception`` (уйдёт в GlitchTip), а в HTTP отдаём generic detail
с тем же error-id для корреляции.
ВАЖНО: вызывать только внутри ``except``-блока — ``logger.exception``
полагается на активный exc_info.
"""
error_id = uuid4().hex[:8]
logger.exception("%s (error-id=%s)", log_context, error_id)
return HTTPException(status_code=status_code, detail=f"{public_detail} (error-id: {error_id})")
_ALL_SOURCES = ["avito", "cian", "yandex"]
class ScrapeRequest(BaseModel):
lat: float
lon: float
radius_m: int = Field(default=1000, ge=100, le=20000)
sources: list[Literal["avito", "cian", "yandex"]] = Field(
default_factory=lambda: list(_ALL_SOURCES)
)
# Если True — скрейп Yandex по 5 сегментам комнат, не одним общим запросом.
multi_room_yandex: bool = False
# Если True — Yandex полный обход 5 rooms × 3 sorts × 2 pages = 30 запросов / ~150s.
deep_yandex: bool = False
# Если True — скрейп Cian по 4 сегментам комнат отдельно (~4× лотов).
multi_room_cian: bool = False
class ScrapeResult(BaseModel):
source: str
fetched: int
inserted: int
updated: int
class ScrapeResponse(BaseModel):
total_fetched: int
total_inserted: int
total_updated: int
by_source: list[ScrapeResult]
@router.post("/scrape", response_model=ScrapeResponse)
async def scrape_around(
payload: ScrapeRequest,
db: Annotated[Session, Depends(get_db)],
) -> ScrapeResponse:
"""Запустить парсеры для точки (lat, lon) в радиусе radius_m метров.
Примеры:
curl -X POST /api/v1/admin/scrape \\
-H 'Content-Type: application/json' \\
-d '{"lat":56.8332,"lon":60.5944,"radius_m":1000,"sources":["avito"]}'
"""
results: list[ScrapeResult] = []
# scraper_kit DI (#2305): config/delay_provider/proxy_provider/matcher инжектируются
# снаружи вместо прямого импорта app.core.config.settings внутри скрапперов —
# тот же паттерн, что и в app.scheduler_main._run_kit_scheduler.
config = RealScraperConfig()
proxy_provider = _kit_proxy_provider()
matcher = RealMatcherAdapter()
for source in payload.sources:
scraper_ctx: AvitoScraper | CianScraper | YandexRealtyScraper
if source == "avito":
scraper_ctx = AvitoScraper(
config, delay_provider=get_scraper_delay, proxy_provider=proxy_provider
)
elif source == "cian":
scraper_ctx = CianScraper(
config, delay_provider=get_scraper_delay, proxy_provider=proxy_provider
)
elif source == "yandex":
scraper_ctx = YandexRealtyScraper(
config, delay_provider=get_scraper_delay, proxy_provider=proxy_provider
)
else:
continue
async with scraper_ctx as scraper:
if source == "yandex" and payload.deep_yandex:
lots = await scraper.fetch_around_multi_room(
payload.lat,
payload.lon,
payload.radius_m,
sorts=("DATE_DESC", "PRICE", "AREA_DESC"),
pages=(0, 1),
)
elif source == "yandex" and payload.multi_room_yandex:
lots = await scraper.fetch_around_multi_room(
payload.lat, payload.lon, payload.radius_m
)
elif source == "cian" and payload.multi_room_cian:
lots = await scraper.fetch_around_multi_room(
payload.lat, payload.lon, payload.radius_m
)
else:
lots = await scraper.fetch_around(payload.lat, payload.lon, payload.radius_m)
# run_id нет и не будет (#2701): ручной admin-скрейп строки в scrape_runs не
# заводит — снимок пишется вне прогона, поле честно остаётся NULL.
inserted, updated = save_listings(
db, lots, matcher=matcher, region_code=DEFAULT_REGION_CODE
)
results.append(
ScrapeResult(source=source, fetched=len(lots), inserted=inserted, updated=updated)
)
return ScrapeResponse(
total_fetched=sum(r.fetched for r in results),
total_inserted=sum(r.inserted for r in results),
total_updated=sum(r.updated for r in results),
by_source=results,
)
def _clean_address_for_geocode(addr: str) -> str:
"""Чистим address для геокодера.
Cian отдаёт «улица Латвийская, 56/3 · р-н Чкаловский» — суффикс ' · ...'
мешает Nominatim. Берём часть до ' · '. Остальные источники такого суффикса
не используют — адрес остаётся без изменений.
"""
main = addr.split(" · ")[0].strip()
return main or addr
@router.post("/geocode-missing")
async def geocode_missing(
db: Annotated[Session, Depends(get_db)],
limit: int = 100,
target: Literal["listings", "deals"] = "listings",
) -> dict:
"""Геокодинг listings ИЛИ deals у которых нет lat/lon (используя address).
target=listings (по умолч.) — объявления; target=deals — сделки Росреестра.
Чанк-обработка с бюджетом по времени (~240с, заведомо меньше cron
`curl -m 320`): за вызов геокодим сколько успеваем, остаток уходит в
`remaining`, cron вызывает в цикле пока `remaining` > 0.
geocode_tried_at: после КАЖДОЙ попытки (успех/провал) ставим NOW(). Failed-
адреса не выбираются повторно 7 дней → cron-loop завершается, не зацикливается.
geom обновляется автоматически триггером.
"""
# Доп. фильтр для listings — у Avito встречаются плейсхолдер-адреса.
extra_filter = "AND address NOT LIKE '%(Avito)%'" if target == "listings" else ""
rows = (
db.execute(
text(
f"""
SELECT id, address, city
FROM {target}
WHERE lat IS NULL
AND COALESCE(address, '') != ''
{extra_filter}
AND (geocode_tried_at IS NULL
OR geocode_tried_at < NOW() - interval '7 days')
ORDER BY geocode_tried_at NULLS FIRST
LIMIT :limit
"""
),
{"limit": limit},
)
.mappings()
.all()
)
# Бюджет по времени: геокодинг упирается в Nominatim (1 req/sec + typo-тиры
# на провалах), 100 адресов могут не уложиться в cron-таймаут `curl -m 320`.
# Выходим из цикла заведомо раньше — remaining > 0 заставит cron продолжить.
budget_sec = 240.0
started = time.monotonic()
geocoded = 0
skipped = 0
for row in rows:
if time.monotonic() - started > budget_sec:
logger.info(
"geocode-missing: бюджет %.0fс исчерпан, обработано %d — cron продолжит",
budget_sec,
geocoded + skipped,
)
break
clean = _clean_address_for_geocode(row["address"])
# city (#2594 шаг 2/3) — известен вызывающему коду через listings.city
# (миграция 196) / deals.city (миграция 177), проставляется из контекста
# развёртки/импорта. Прокидываем как city_hint, а не полагаемся на то, что
# геокодер угадает город по тексту address (голый "ул. Победы, 30" без
# города в тексте иначе уходит в Екатеринбург).
#
# known_city_hint (#2603) — гейт по словарю городов области: при
# target="deals" сюда приходит росреестровое поле, в хвосте которого
# лежат не-города («Бессонова», «Билейский рыбопитомник»), а мусорный
# хинт закрывает EKB-локальные тиры и уезжает префиксом в запрос
# провайдеру, т.е. вреднее отсутствия хинта. Общий хелпер, тот же, что у
# scripts/geocode_deals_nominatim.py и tasks/geocode_missing.py.
city = known_city_hint(row.get("city"))
result = await geocode(clean, db, city_hint=city)
if result is None:
# Помечаем что пробовали — иначе ретрай на каждом cron.
db.execute(
text(f"UPDATE {target} SET geocode_tried_at = NOW() WHERE id = :id"),
{"id": row["id"]},
)
db.commit()
skipped += 1
continue
# geom пересчитается триггером из lat/lon.
db.execute(
text(
f"UPDATE {target} SET lat = :lat, lon = :lon, "
"geocode_tried_at = NOW() WHERE id = :id"
),
{"lat": result.lat, "lon": result.lon, "id": row["id"]},
)
db.commit()
geocoded += 1
# Сколько ещё не-геокоженных адресов ждут (для cron-loop'а).
remaining = db.execute(
text(
f"""
SELECT count(*) FROM {target}
WHERE lat IS NULL
AND COALESCE(address, '') != ''
{extra_filter}
AND (geocode_tried_at IS NULL
OR geocode_tried_at < NOW() - interval '7 days')
"""
)
).scalar()
return {
"checked": len(rows),
"geocoded": geocoded,
"skipped": skipped,
"remaining": int(remaining or 0),
}
@router.post("/scrape/cian/upload-cookies", status_code=200)
async def upload_cian_cookies(
cookies: dict[str, str],
db: Annotated[Session, Depends(get_db)],
) -> dict:
"""Upload Cian session cookies для аутентификации Valuation Calculator.
Body: JSON object вида { "_ym_uid": "...", "_cian_visitor_session_id": "...", ... }
Cookies экспортируются из Chrome DevTools → Application → Cookies → www.cian.ru
или через расширение Cookie-Editor → Export → JSON (object format).
Шаги:
1. Фильтрует payload по CIAN_REQUIRED_COOKIES.
2. Верифицирует через GET /kalkulator-nedvizhimosti/ — проверяет isAuthenticated.
3. Сохраняет зашифрованно (pgp_sym_encrypt) в cian_session_cookies.
Returns: {"ok": true, "userId": <int>, "cookieCount": <int>}
"""
if not settings.cookie_encryption_key:
raise HTTPException(status_code=503, detail="COOKIE_ENCRYPTION_KEY not configured")
cleaned = {k: v for k, v in cookies.items() if k in cian_session_svc.CIAN_REQUIRED_COOKIES}
if not cleaned:
raise HTTPException(
status_code=400,
detail=(
f"No recognized Cian cookies in payload. "
f"Expected any of: {sorted(cian_session_svc.CIAN_REQUIRED_COOKIES)}"
),
)
state = await cian_session_svc.verify_session(cleaned)
if state is None:
raise HTTPException(
status_code=401,
detail="Cookies invalid or session not authenticated on cian.ru",
)
user_id = state.get("user", {}).get("userId")
if not user_id:
raise HTTPException(status_code=400, detail="Authenticated state missing userId")
cian_session_svc.save_session(db, account_user_id=int(user_id), cookies=cleaned)
return {"ok": True, "userId": user_id, "cookieCount": len(cleaned)}
class CianAutoLoginRequest(BaseModel):
email: str | None = None # override settings.cian_login_email
password: str | None = None # override settings.cian_login_password
@router.post("/scrape/cian/auto-login", status_code=200)
async def cian_auto_login(
db: Annotated[Session, Depends(get_db)],
body: CianAutoLoginRequest | None = None,
) -> dict:
"""Browser auto-login на cian.ru (Variant B, #639) → извлекает cookies → save_session.
Headless camoufox (tradein-browser /login) логинится email+паролем, забирает
auth-cookies, фильтрует по CIAN_REQUIRED_COOKIES, верифицирует isAuthenticated,
сохраняет зашифрованно. Креды из settings.cian_login_* или из body (override).
Returns: {"ok": true, "userId": <int>, "cookieCount": <int>}
"""
if not settings.cookie_encryption_key:
raise HTTPException(status_code=503, detail="COOKIE_ENCRYPTION_KEY not configured")
email = (body.email if body else None) or settings.cian_login_email
password = (body.password if body else None) or settings.cian_login_password
if not email or not password:
raise HTTPException(
status_code=503,
detail="CIAN_LOGIN_EMAIL/PASSWORD not configured",
)
try:
async with BrowserFetcher(
source="cian", endpoint=settings.browser_http_endpoint
) as fetcher:
raw_cookies = await fetcher.login(
url=settings.cian_login_url,
email=email,
password=password,
email_selector=settings.cian_login_email_selector,
password_selector=settings.cian_login_password_selector,
submit_selector=settings.cian_login_submit_selector,
success_cookie=settings.cian_login_success_cookie,
pre_click_selectors=settings.cian_login_pre_click_selectors,
wait_ms=settings.cian_login_wait_ms,
)
except Exception as exc:
logger.error("cian auto-login failed: %s", type(exc).__name__)
raise HTTPException(
status_code=502, detail="Browser login failed (see server logs)"
) from exc
cleaned = {k: v for k, v in raw_cookies.items() if k in cian_session_svc.CIAN_REQUIRED_COOKIES}
if not cleaned:
raise HTTPException(
status_code=502,
detail=(
"Login returned no recognized Cian cookies "
"(selectors may be stale — check CIAN_LOGIN_*_SELECTOR)"
),
)
state = await cian_session_svc.verify_session(cleaned)
if state is None:
raise HTTPException(
status_code=401,
detail="Logged in but session not authenticated (cookies rejected by cian.ru)",
)
user_id = state.get("user", {}).get("userId")
if not user_id:
raise HTTPException(status_code=400, detail="Authenticated state missing userId")
cian_session_svc.save_session(db, account_user_id=int(user_id), cookies=cleaned)
return {"ok": True, "userId": user_id, "cookieCount": len(cleaned)}
@router.get("/scrape/cian/test-auth", status_code=200)
async def test_cian_auth(
db: Annotated[Session, Depends(get_db)],
) -> dict:
"""Проверить что текущие сохранённые Cian cookies ещё валидны.
Returns: {"authenticated": bool, "userId": <int|null>, "reason": <str|null>}
"""
if not settings.cookie_encryption_key:
return {"authenticated": False, "userId": None, "reason": "encryption_key_not_configured"}
cookies = cian_session_svc.load_session(db)
if cookies is None:
return {"authenticated": False, "userId": None, "reason": "no_session_in_db"}
state = await cian_session_svc.verify_session(cookies)
if state is None:
return {"authenticated": False, "userId": None, "reason": "session_expired_or_invalid"}
if state.get("_ban"):
return {"authenticated": False, "userId": None, "reason": "banned_403"}
user_id = state.get("user", {}).get("userId")
return {"authenticated": True, "userId": user_id, "reason": None}
# ── DomClick session cookie management (cookie-injection MVP, #2000) ─────────
#
# Эмпирически подтверждено вживую (2026-07-04): QRATOR-блок DomClick (даже на
# уже подозрительном прокси-IP) полностью обходится инъекцией cookies валидной
# аутентифицированной test-аккаунт сессии (Sber ID) перед навигацией на карточку.
# В отличие от Cian, здесь нет verify_session (потребовал бы реального browser-
# фетча, не curl_cffi) — MVP просто фильтрует+сохраняет, доверяя caller'у.
@router.post("/scrape/domclick/upload-cookies", status_code=200)
async def upload_domclick_cookies(
cookies: dict[str, str],
db: Annotated[Session, Depends(get_db)],
) -> dict:
"""Upload DomClick session cookies (Sber ID) для обхода QRATOR anti-bot блока.
Body: JSON object вида { "CAS_ID": "...", "qrator_jsid2": "...", ... } — плоский
name→value dict (caller уже флэттенит сырой Cookie-Editor export до этой формы;
полный JSON-массив формат Cookie-Editor здесь НЕ парсим).
Cookies экспортируются из Chrome DevTools → Application → Cookies →
*.domclick.ru или через расширение Cookie-Editor → Export → JSON.
Шаги:
1. Фильтрует payload по DOMCLICK_REQUIRED_COOKIES.
2. Извлекает account_cas_id из cookies["CAS_ID"] (int).
3. Сохраняет зашифрованно (pgp_sym_encrypt) в domclick_session_cookies.
MVP: без verify — в отличие от Cian, DomClick verification требует реального
browser-фетча (future enhancement), не простого curl_cffi-запроса.
Returns: {"ok": true, "accountCasId": <int>, "cookieCount": <int>}
"""
if not settings.cookie_encryption_key:
raise HTTPException(status_code=503, detail="COOKIE_ENCRYPTION_KEY not configured")
cleaned = {
k: v for k, v in cookies.items() if k in domclick_session_svc.DOMCLICK_REQUIRED_COOKIES
}
if not cleaned:
raise HTTPException(
status_code=400,
detail=(
f"No recognized DomClick cookies in payload. "
f"Expected any of: {sorted(domclick_session_svc.DOMCLICK_REQUIRED_COOKIES)}"
),
)
raw_cas_id = cookies.get("CAS_ID")
if raw_cas_id is None:
raise HTTPException(
status_code=400, detail="Missing CAS_ID cookie — cannot derive account id"
)
try:
account_cas_id = int(raw_cas_id)
except (TypeError, ValueError):
raise HTTPException(status_code=400, detail="CAS_ID cookie is not numeric") from None
domclick_session_svc.save_session(db, account_cas_id=account_cas_id, cookies=cleaned)
return {"ok": True, "accountCasId": account_cas_id, "cookieCount": len(cleaned)}
class DomClickDebugDetailFetchRequest(BaseModel):
card_url: str
class DomClickDebugDetailFetchResponse(BaseModel):
ok: bool
card_url: str
cookies_injected: bool
item_id: str | None
repair_state: str | None
living_area_m2: float | None
year_built: int | None
price_changes_count: int
raw_extra: dict[str, Any]
@router.post("/scrape/domclick/debug/detail-fetch", response_model=DomClickDebugDetailFetchResponse)
async def debug_domclick_detail_fetch(
body: DomClickDebugDetailFetchRequest,
db: Annotated[Session, Depends(get_db)],
) -> DomClickDebugDetailFetchResponse:
"""Ad-hoc fetch одной DomClick detail-карточки с cookie-инъекцией (debug, #2000).
Единственный СЕЙЧАС живой путь прогнать связку domclick_session.load_session +
BrowserFetcher(cookies=...) + scraper_kit.providers.domclick.detail.fetch_detail —
боевой backfill-оркестратор для DomClick ещё не смёржен (issue #2000).
cookies=None (валидной session нет в БД) — fetch без инъекции, остаётся только
origin-warmup через SERP (#2430). Debug-only, ничего не пишет в БД.
"""
_assert_domclick_host(body.card_url)
# Local import — зеркалит паттерн scrape_cian_detail/scrape_cian_newbuilding в
# этом же файле: domclick.detail.fetch_detail тёзка module-level avito fetch_detail
# (строка 51), local import ограничивает shadowing только этой функцией.
from scraper_kit.providers.domclick.detail import fetch_detail as domclick_fetch_detail
cookies = domclick_session_svc.load_session(db)
# source="domclick" — выделенный резидентный прокси (scrape_proxies.id=1,
# provider_affinity='domclick', см. 173_scrape_proxies_add_domclick_affinity.sql)
# построен и verified live именно 2026-07-04 (815КБ __SSR_STATE__ через свежий IP).
# Заменяет прежний source="cian" (docstring domclick/detail.py — устаревшая
# рекомендация от 2026-06-27, до появления выделенного пула).
async with BrowserFetcher(source="domclick", endpoint=settings.browser_http_endpoint) as bf:
try:
enrichment = await domclick_fetch_detail(
body.card_url, browser_fetcher=bf, cookies=cookies
)
except Exception as e:
raise _safe_http_error(
502,
"fetch failed",
f"domclick-debug-detail-fetch: fetch failed for {body.card_url}",
) from e
return DomClickDebugDetailFetchResponse(
ok=True,
card_url=body.card_url,
cookies_injected=cookies is not None,
item_id=enrichment.item_id,
repair_state=enrichment.repair_state,
living_area_m2=enrichment.living_area_m2,
year_built=enrichment.year_built,
price_changes_count=len(enrichment.price_changes),
raw_extra=enrichment.raw_extra,
)
# ── Geocode backfill: batch address-dedup geocoder (all sources) ─────────────
@router.get("/scrape/geocode-missing-listings/status")
async def get_geocode_backfill_status(
db: Annotated[Session, Depends(get_db)],
) -> dict[str, Any]:
"""Текущая статистика backfill — сколько ещё pending geocode.
Включает per-source breakdown и оценку времени при Nominatim 1 req/sec.
"""
stats = (
db.execute(
text(
"""
SELECT
source,
COUNT(*) AS total,
COUNT(*) FILTER (WHERE lat IS NULL AND address IS NOT NULL) AS pending_geocode,
COUNT(*) FILTER (WHERE lat IS NOT NULL) AS has_coords
FROM listings
WHERE is_active = true
GROUP BY source
ORDER BY pending_geocode DESC
"""
)
)
.mappings()
.all()
)
total_pending = sum(s["pending_geocode"] for s in stats)
unique_addresses_pending = (
db.execute(
text(
"""
SELECT COUNT(DISTINCT address) FROM listings
WHERE lat IS NULL AND address IS NOT NULL AND length(trim(address)) >= 5
"""
)
).scalar()
or 0
)
cache_size = (
db.execute(text("SELECT COUNT(*) FROM geocode_cache WHERE expires_at > NOW()")).scalar()
or 0
)
return {
"total_pending_listings": total_pending,
"unique_addresses_pending": int(unique_addresses_pending),
"estimated_nominatim_minutes": round(int(unique_addresses_pending) / 60.0, 1),
"cache_size_active": int(cache_size),
"per_source": [dict(s) for s in stats],
}
@router.post("/scrape/geocode-missing-listings")
async def trigger_geocode_missing_listings(
background_tasks: BackgroundTasks,
db: Annotated[Session, Depends(get_db)],
batch_size: int = Query(default=200, ge=1, le=2000),
dry_run: bool = Query(default=False),
background: bool = Query(
default=False,
description="True → запустить в BackgroundTasks (не ждать результат).",
),
) -> dict[str, Any]:
"""Manual trigger для geocode backfill всех listings с NULL coords (address-dedup).
batch_size=200 ≈ 3.5 мин при Nominatim 1 req/sec (без cache hits — быстрее с cache).
Безопасно для повторного запуска (ON CONFLICT в geocode_cache, UPDATE WHERE lat IS NULL).
dry_run=true — показывает что бы сделалось, без UPDATE listings.
background=true — запускает через BackgroundTasks, возвращает немедленно (для больших batch).
"""
from app.core.db import SessionLocal
if background:
async def _bg_task() -> None:
bg_db = SessionLocal()
try:
await geocode_missing_listings(bg_db, batch_size=batch_size, dry_run=dry_run)
except Exception:
logger.exception("geocode-missing-listings background task crashed")
finally:
bg_db.close()
background_tasks.add_task(_bg_task)
return {
"status": "queued",
"batch_size": batch_size,
"dry_run": dry_run,
"note": "Running in background. Check logs for progress.",
}
result: GeocodeBackfillResult = await geocode_missing_listings(
db, batch_size=batch_size, dry_run=dry_run
)
return {
"status": "ok",
"batch_size": batch_size,
"dry_run": dry_run,
"addresses_total": result.addresses_total,
"addresses_processed": result.addresses_processed,
"addresses_geocoded": result.addresses_geocoded,
"addresses_failed": result.addresses_failed,
"listings_updated": result.listings_updated,
"cache_hits": result.cache_hits,
"cache_misses": result.cache_misses,
"duration_sec": round(result.duration_sec, 2),
}
# ── Stage 4a: Avito enrichment endpoints (single-shot debug + manual triggers) ──
@router.post("/scrape/avito-house")
async def scrape_avito_house(
db: Annotated[Session, Depends(get_db)],
house_url: str,
) -> dict[str, object]:
"""Enrichment ОДНОГО дома Avito Houses Catalog.
Body params (query): `house_url=/catalog/houses/<slug>/<id>` (absolute или relative path).
Flow: fetch_house_catalog → save_house_catalog_enrichment.
Returns counters {'house_id', 'reviews', 'sellers', 'listings_linked', 'placement_history'}.
"""
_assert_allowed_url(house_url)
try:
enrichment = await fetch_house_catalog(house_url)
except Exception as e:
raise _safe_http_error(
502, "fetch failed", f"avito-house: fetch failed for {house_url}"
) from e
try:
counters = save_house_catalog_enrichment(db, enrichment, matcher=RealMatcherAdapter())
except Exception as e:
raise _safe_http_error(
500, "save failed", f"avito-house: save failed for {house_url}"
) from e
logger.info("avito-house ok: %s%s", house_url, counters)
return {"ok": True, "house_url": house_url, "counters": counters}
@router.post("/scrape/avito-detail")
async def scrape_avito_detail(
db: Annotated[Session, Depends(get_db)],
item_url: str,
) -> dict[str, object]:
"""Detail enrichment ОДНОГО listing.
Body params (query): `item_url=/ekaterinburg/kvartiry/...` (absolute или relative).
Flow: fetch_detail → save_detail_enrichment (UPDATE listings WHERE source='avito').
Returns {'ok', 'item_id', 'updated'}.
"""
_assert_allowed_url(item_url)
try:
enrichment = await fetch_detail(item_url, config=RealScraperConfig())
except Exception as e:
raise _safe_http_error(
502, "fetch failed", f"avito-detail: fetch failed for {item_url}"
) from e
try:
updated = save_detail_enrichment(db, enrichment)
except Exception as e:
raise _safe_http_error(
500, "save failed", f"avito-detail: save failed for {item_url}"
) from e
if not updated:
# Listing с этим source_id ещё не существует в БД — был bypass через
# /admin/scrape/avito-detail для несуществующего listing.
logger.warning("avito-detail: no listing matched source_id=%s", enrichment.item_id)
logger.info("avito-detail ok: %s item_id=%s updated=%s", item_url, enrichment.item_id, updated)
return {"ok": True, "item_id": enrichment.item_id, "updated": updated}
@router.post("/scrape/avito-detail-backfill")
async def scrape_avito_detail_backfill(
db: Annotated[Session, Depends(get_db)],
batch_size: int = Query(default=200, ge=1, le=2000),
budget_sec: float = Query(default=600.0, ge=30.0, le=3600.0),
) -> dict[str, object]:
"""On-demand прогон detail-backfill detail-очереди avito (#1950).
Запускает `run_avito_detail_backfill` синхронно (await) на одном snapshot'е
из batch_size pending-листингов (detail_enriched_at IS NULL, активные ЕКБ),
ограниченный budget_sec. Создаёт scrape_runs-строку (source='avito_detail_backfill')
для observability — статус виден в общем списке runs.
Нужен для ручной верификации/прогона: основной schedule window-gated (09-12 UTC),
вне окна его не запустить. Не для постоянного использования — backlog нормально
рассасывается ночным schedule.
Returns counters {'attempted','enriched','blocked','failed','duration_sec'} + run_id.
"""
if has_running_run(db, "avito_detail_backfill"):
raise HTTPException(
status_code=409,
detail=(
"avito detail-backfill уже выполняется — одновременно допустим только один "
"(параллельные раны делят прокси-IP → Avito бан). Дождитесь завершения."
),
)
params = {"batch_size": batch_size, "budget_sec": budget_sec}
run_id = runs_mod.create_run(db, source="avito_detail_backfill", params=params)
try:
result = await run_avito_detail_backfill(db, run_id=run_id, params=params)
except Exception as e:
raise _safe_http_error(
502, "backfill failed", f"avito-detail-backfill: run_id={run_id} crashed"
) from e
logger.info("avito-detail-backfill ok: run_id=%d %s", run_id, result.to_dict())
return {"ok": True, "run_id": run_id, "counters": result.to_dict()}
@router.post("/scrape/avito-imv")
async def scrape_avito_imv(
db: Annotated[Session, Depends(get_db)],
address: str,
rooms: int,
area_m2: float,
floor: int,
floor_at_home: int,
house_type: str,
renovation_type: str,
has_balcony: bool = False,
has_loggia: bool = False,
) -> dict[str, object]:
"""Debug-trigger Avito IMV evaluation. БЕЗ cache — всегда свежий fetch.
Production IMV вызывается из estimator (Stage 3, on-demand с 24h cache).
Этот endpoint — для ручной отладки / data-pipeline tools.
Returns {'ok', 'cache_key', 'recommended_price', 'range', 'market_count',
'evaluation_id', 'history_saved'}.
"""
try:
result = await evaluate_via_imv(
address=address,
rooms=rooms,
area_m2=area_m2,
floor=floor,
floor_at_home=floor_at_home,
house_type=house_type,
renovation_type=renovation_type,
has_balcony=has_balcony,
has_loggia=has_loggia,
config=RealScraperConfig(),
)
except IMVAddressNotFoundError as e:
# Ожидаемое клиентское условие (адрес не в базе Avito), НЕ сбой — logger.warning
# без traceback, чтобы не шуметь exception-событиями в GlitchTip. Адрес — в лог,
# не в HTTP-ответ.
error_id = uuid4().hex[:8]
logger.warning("avito-imv: address not found for %s (error-id=%s)", address, error_id)
raise HTTPException(
status_code=404, detail=f"IMV address not found (error-id: {error_id})"
) from e
except Exception as e:
raise _safe_http_error(502, "fetch failed", f"avito-imv: fetch failed for {address}") from e
try:
eval_id = save_imv_evaluation(db, result)
history_saved = save_imv_placement_history(db, eval_id, result.placement_history)
except Exception as e:
raise _safe_http_error(500, "save failed", f"avito-imv: save failed for {address}") from e
logger.info(
"avito-imv ok: addr=%s recommended=%d evaluation_id=%d history=%d",
address[:60],
result.recommended_price,
eval_id,
history_saved,
)
return {
"ok": True,
"cache_key": result.cache_key,
"recommended_price": result.recommended_price,
"range": [result.lower_price, result.higher_price],
"market_count": result.market_count,
"evaluation_id": eval_id,
"history_saved": history_saved,
}
# ── City sweep: background full-ЕКБ pipeline ────────────────────────────────
class CitySweepStartRequest(BaseModel):
pages_per_anchor: int = Field(default=3, ge=1, le=10)
detail_top_n: int = Field(default=10, ge=0, le=30)
request_delay_sec: float = Field(default=7.0, ge=3.0, le=15.0)
enrich_houses: bool = True
enrich_imv: bool = Field(
default=True,
description=(
"Если True — после всех anchor'ов запускает IMV-оценку тронутых домов "
"(финальная фаза sweep'а). Требует enrich_houses=True для сбора house_id."
),
)
radius_m: int = Field(default=1500, ge=500, le=5000)
class CitySweepStartResponse(BaseModel):
run_id: int
status: str
pages_per_anchor: int
detail_top_n: int
class ScrapeRunRow(BaseModel):
run_id: int
source: str
status: str
params: dict | None = None
counters: dict | None = None
error: str | None = None
started_at: str | None = None
finished_at: str | None = None
heartbeat_at: str | None = None
@router.post("/scrape/avito-city-sweep", response_model=CitySweepStartResponse)
async def start_avito_city_sweep(
payload: CitySweepStartRequest,
background_tasks: BackgroundTasks,
db: Annotated[Session, Depends(get_db)],
) -> CitySweepStartResponse:
"""Запустить full ЕКБ city sweep в background. Returns run_id для polling/cancel.
Sweep итерирует 5 anchor-точек ЕКБ × pages_per_anchor страниц Avito,
сохраняет listings, обогащает houses и detail. Runs in FastAPI BackgroundTasks.
Coop cancel: POST /scrape/avito-city-sweep/{run_id}/cancel помечает run,
pipeline проверяет статус каждый anchor.
Single-run guard: второй sweep запустить нельзя, пока первый running —
параллельные раны делят один прокси-IP → удвоенный rate → Avito бан
(incident 2026-05-31: runs #26+#27 шли одновременно → оба banned).
Scheduler уже защищён (has_running_run + advisory lock); этот guard
закрывает ручной admin-триггер.
"""
if has_running_run(db, "avito_city_sweep"):
raise HTTPException(
status_code=409,
detail=(
"avito city-sweep уже выполняется — одновременно допустим только один "
"(параллельные раны делят прокси-IP → Avito бан). Дождитесь завершения "
"или отмените текущий через POST /scrape/avito-city-sweep/{run_id}/cancel."
),
)
run_id = runs_mod.create_run(db, source="avito_city_sweep", params=payload.model_dump())
async def _sweep_task() -> None:
sweep_db = SessionLocal()
try:
await run_avito_city_sweep(
sweep_db,
run_id=run_id,
config=RealScraperConfig(),
matcher=RealMatcherAdapter(),
enrichment=RealEnrichmentJobs(),
proxy_provider=_kit_proxy_provider(),
radius_m=payload.radius_m,
pages_per_anchor=payload.pages_per_anchor,
enrich_houses=payload.enrich_houses,
detail_top_n=payload.detail_top_n,
request_delay_sec=payload.request_delay_sec,
enrich_imv=payload.enrich_imv,
)
except Exception:
logger.exception("city-sweep background task run_id=%d crashed", run_id)
finally:
sweep_db.close()
background_tasks.add_task(_sweep_task)
logger.info("city-sweep queued run_id=%d params=%s", run_id, payload.model_dump())
return CitySweepStartResponse(
run_id=run_id,
status="running",
pages_per_anchor=payload.pages_per_anchor,
detail_top_n=payload.detail_top_n,
)
@router.get("/scrape/avito-city-sweep/runs", response_model=list[ScrapeRunRow])
def list_avito_city_sweep_runs(
db: Annotated[Session, Depends(get_db)],
limit: int = 10,
) -> list[ScrapeRunRow]:
"""Список последних N runs (для UI polling). Default limit=10."""
rows = runs_mod.list_recent(db, source="avito_city_sweep", limit=limit)
return [
ScrapeRunRow(
run_id=r["run_id"],
source=r["source"],
status=r["status"],
params=r.get("params"),
counters=r.get("counters"),
error=r.get("error"),
started_at=r["started_at"].isoformat() if r.get("started_at") else None,
finished_at=r["finished_at"].isoformat() if r.get("finished_at") else None,
heartbeat_at=r["heartbeat_at"].isoformat() if r.get("heartbeat_at") else None,
)
for r in rows
]
@router.post("/scrape/avito-city-sweep/{run_id}/cancel")
def cancel_avito_city_sweep(
run_id: int,
db: Annotated[Session, Depends(get_db)],
) -> dict[str, object]:
"""Отменить running sweep. Cooperative: pipeline проверяет scrape_runs.status каждый anchor."""
cancelled = runs_mod.mark_cancelled(db, run_id)
return {"ok": True, "run_id": run_id, "cancelled": cancelled}
# ── Cian city sweep (#860) — on-demand trigger + runs list ───────────────────
class CianCitySweepStartRequest(BaseModel):
pages_per_anchor: int = Field(default=3, ge=1, le=10)
detail_top_n: int = Field(default=10, ge=0, le=30)
request_delay_sec: float = Field(default=5.0, ge=3.0, le=15.0)
enrich_houses: bool = True
radius_m: int = Field(default=1500, ge=500, le=5000)
@router.post("/scrape/cian-city-sweep", response_model=CitySweepStartResponse)
async def start_cian_city_sweep(
payload: CianCitySweepStartRequest,
background_tasks: BackgroundTasks,
db: Annotated[Session, Depends(get_db)],
) -> CitySweepStartResponse:
"""Запустить Cian ЕКБ city sweep в background (#860). Returns run_id.
Один прогон: SERP (fetch_around_multi_room) → detail-обогащение (incl.
price-history) → newbuilding/houses (если enrich_houses=True).
Прокси уже внутри CianScraper.
Coop cancel: POST /scrape/cian-city-sweep/{run_id}/cancel.
НЕ добавляется в scrape_schedules — ручной/контролируемый запуск
(bulk cron → 429-бан на Cian).
"""
run_id = runs_mod.create_run(db, source="cian_city_sweep", params=payload.model_dump())
async def _sweep_task() -> None:
sweep_db = SessionLocal()
try:
await run_cian_city_sweep(
sweep_db,
run_id=run_id,
config=RealScraperConfig(),
matcher=RealMatcherAdapter(),
proxy_provider=_kit_proxy_provider(),
radius_m=payload.radius_m,
pages_per_anchor=payload.pages_per_anchor,
request_delay_sec=payload.request_delay_sec,
detail_top_n=payload.detail_top_n,
enrich_houses=payload.enrich_houses,
)
except Exception:
logger.exception("cian-sweep background task run_id=%d crashed", run_id)
finally:
sweep_db.close()
background_tasks.add_task(_sweep_task)
logger.info("cian-sweep queued run_id=%d params=%s", run_id, payload.model_dump())
return CitySweepStartResponse(
run_id=run_id,
status="running",
pages_per_anchor=payload.pages_per_anchor,
detail_top_n=payload.detail_top_n,
)
@router.get("/scrape/cian-city-sweep/runs", response_model=list[ScrapeRunRow])
def list_cian_city_sweep_runs(
db: Annotated[Session, Depends(get_db)],
limit: int = 10,
) -> list[ScrapeRunRow]:
"""Список последних N Cian city sweep runs (для UI polling). Default limit=10."""
rows = runs_mod.list_recent(db, source="cian_city_sweep", limit=limit)
return [
ScrapeRunRow(
run_id=r["run_id"],
source=r["source"],
status=r["status"],
params=r.get("params"),
counters=r.get("counters"),
error=r.get("error"),
started_at=r["started_at"].isoformat() if r.get("started_at") else None,
finished_at=r["finished_at"].isoformat() if r.get("finished_at") else None,
heartbeat_at=r["heartbeat_at"].isoformat() if r.get("heartbeat_at") else None,
)
for r in rows
]
@router.post("/scrape/cian-city-sweep/{run_id}/cancel")
def cancel_cian_city_sweep(
run_id: int,
db: Annotated[Session, Depends(get_db)],
) -> dict[str, object]:
"""Отменить running Cian sweep. Cooperative: проверяется каждый anchor."""
cancelled = runs_mod.mark_cancelled(db, run_id)
return {"ok": True, "run_id": run_id, "cancelled": cancelled}
# ── Cian exhaustive full load (room × price partitioning) ────────────────────
class CianFullLoadRequest(BaseModel):
price_cap_per_bucket: int = Field(
default=1400,
ge=200,
le=1500,
description=(
"Максимум офферов на price-бакет. При totalOffers > cap — бакет делится пополам. "
"Ниже Cian SERP-cap ~1500; запас 100 на variance."
),
)
request_delay_sec: float = Field(
default=1.5,
ge=0.5,
le=15.0,
description=(
"Задержка между SERP-запросами (сек). Выделенный прокси + IP-ротация "
"позволяют снизить до 1.5s. Увеличить при бане."
),
)
concurrency: int = Field(
default=5,
ge=1,
le=10,
description="Параллельных page-фетчей внутри одного leaf-бакета.",
)
enrich_detail: bool = Field(
default=False,
description="Если True — после сбора обогатить detail-страницы (медленно).",
)
detail_top_n: int = Field(
default=0,
ge=0,
le=50,
description="Сколько свежих листингов обогатить detail (актуально при enrich_detail=True).",
)
resume_run_id: int | None = Field(
default=None,
description=(
"Если задан — читает done_buckets из counters прошлого run и пропускает "
"уже завершённые бакеты. Передать run_id предыдущего (прерванного) "
"cian_full_load прогона для resume. Без этого поля — full walk с нуля."
),
)
@router.post("/scrape/cian-full-load", response_model=CitySweepStartResponse)
async def start_cian_full_load(
payload: CianFullLoadRequest,
background_tasks: BackgroundTasks,
db: Annotated[Session, Depends(get_db)],
) -> CitySweepStartResponse:
"""Запустить exhaustive Cian ЕКБ full load в background.
Обходит Cian SERP-cap (~54 стр/запрос ≈ 1500 результатов): адаптивно бьёт
цену на бакеты внутри каждой комнатности (totalOffers-driven бинарное деление),
пагинирует каждый бакет полностью, дедуп по source_id.
Один региональный проход БЕЗ anchor'ов (Cian SERP region-wide — anchors избыточны).
IP-ротация через changeip при timeout/captcha.
Coop cancel: POST /scrape/cian-full-load/{run_id}/cancel.
Статус прогона: смотреть через GET /scrape/cian-city-sweep/runs
(или добавить source='cian_full_load' в /admin SQL).
"""
run_id = runs_mod.create_run(db, source="cian_full_load", params=payload.model_dump())
async def _full_load_task() -> None:
task_db = SessionLocal()
try:
await run_cian_full_load(
task_db,
run_id=run_id,
config=RealScraperConfig(),
matcher=RealMatcherAdapter(),
proxy_provider=_kit_proxy_provider(),
price_cap_per_bucket=payload.price_cap_per_bucket,
request_delay_sec=payload.request_delay_sec,
concurrency=payload.concurrency,
enrich_detail=payload.enrich_detail,
detail_top_n=payload.detail_top_n,
resume_run_id=payload.resume_run_id,
)
except Exception:
logger.exception("cian-full-load background task run_id=%d crashed", run_id)
finally:
task_db.close()
background_tasks.add_task(_full_load_task)
logger.info("cian-full-load queued run_id=%d params=%s", run_id, payload.model_dump())
return CitySweepStartResponse(
run_id=run_id,
status="running",
pages_per_anchor=0, # нет anchor'ов в full load
detail_top_n=payload.detail_top_n,
)
@router.post("/scrape/cian-full-load/{run_id}/cancel")
def cancel_cian_full_load(
run_id: int,
db: Annotated[Session, Depends(get_db)],
) -> dict[str, object]:
"""Отменить running Cian full load. Cooperative: проверяется каждый room-bucket."""
cancelled = runs_mod.mark_cancelled(db, run_id)
return {"ok": True, "run_id": run_id, "cancelled": cancelled}
# ── Yandex exhaustive full load (room × price partitioning) ──────────────────
class YandexFullLoadRequest(BaseModel):
price_cap_per_bucket: int = Field(
default=500,
ge=100,
le=575,
description=(
"Максимум офферов на price-бакет. При totalItems > cap — бакет делится пополам. "
"Ниже Yandex SERP-cap ~575; запас на variance. "
"При недоступности totalItems — fallback на paginate-until-empty."
),
)
request_delay_sec: float = Field(
default=2.0,
ge=1.0,
le=15.0,
description=(
"Задержка между SERP-запросами (сек). "
"Рекомендуется ≥ 2.0 — Yandex агрессивнее Cian по captcha-детекции."
),
)
concurrency: int = Field(
default=4,
ge=1,
le=8,
description="Параллельных page-фетчей внутри одного leaf-бакета.",
)
resume_run_id: int | None = Field(
default=None,
description=(
"Если задан — читает done_buckets из counters прошлого run и пропускает "
"уже завершённые бакеты. Передать run_id предыдущего (прерванного) "
"yandex_full_load прогона для resume. Без этого поля — full walk с нуля."
),
)
@router.post("/scrape/yandex-full-load", response_model=CitySweepStartResponse)
async def start_yandex_full_load(
payload: YandexFullLoadRequest,
background_tasks: BackgroundTasks,
db: Annotated[Session, Depends(get_db)],
) -> CitySweepStartResponse:
"""Запустить exhaustive Yandex ЕКБ full load в background.
Обходит Yandex SERP-cap (~575 карточек/запрос): адаптивно бьёт цену на бакеты
внутри каждой комнатности, пагинирует каждый бакет полностью, дедуп по source_id.
При недоступности totalItems из state — автоматический fallback на paginate-until-empty.
Один региональный проход БЕЗ anchor'ов (Yandex SERP path-based по городу — anchors
избыточны). IP-ротация через changeip при captcha.
Checkpoint/resume: задай resume_run_id=<предыдущий run_id> чтобы пропустить
уже завершённые бакеты и продолжить с места остановки.
Coop cancel: POST /scrape/yandex-full-load/{run_id}/cancel.
Статус прогона: смотреть через GET /scrape/cian-city-sweep/runs
(source='yandex_full_load').
"""
run_id = runs_mod.create_run(db, source="yandex_full_load", params=payload.model_dump())
async def _full_load_task() -> None:
task_db = SessionLocal()
try:
await run_yandex_full_load(
task_db,
run_id=run_id,
config=RealScraperConfig(),
matcher=RealMatcherAdapter(),
enrichment=RealEnrichmentJobs(),
proxy_provider=_kit_proxy_provider(),
price_cap_per_bucket=payload.price_cap_per_bucket,
request_delay_sec=payload.request_delay_sec,
concurrency=payload.concurrency,
resume_run_id=payload.resume_run_id,
)
except Exception:
logger.exception("yandex-full-load background task run_id=%d crashed", run_id)
finally:
task_db.close()
background_tasks.add_task(_full_load_task)
logger.info("yandex-full-load queued run_id=%d params=%s", run_id, payload.model_dump())
return CitySweepStartResponse(
run_id=run_id,
status="running",
pages_per_anchor=0, # нет anchor'ов в full load
detail_top_n=0,
)
@router.post("/scrape/yandex-full-load/{run_id}/cancel")
def cancel_yandex_full_load(
run_id: int,
db: Annotated[Session, Depends(get_db)],
) -> dict[str, object]:
"""Отменить running Yandex full load. Cooperative: проверяется каждый room-bucket."""
cancelled = runs_mod.mark_cancelled(db, run_id)
return {"ok": True, "run_id": run_id, "cancelled": cancelled}
# ── Yandex city sweep (#861) — on-demand trigger + runs list ─────────────────
class YandexCitySweepStartRequest(BaseModel):
pages_per_anchor: int = Field(default=2, ge=1, le=10)
request_delay_sec: float = Field(default=7.0, ge=3.0, le=30.0)
enrich_address: bool = True
radius_m: int = Field(default=1500, ge=500, le=5000)
@router.post("/scrape/yandex-city-sweep", response_model=CitySweepStartResponse)
async def start_yandex_city_sweep(
payload: YandexCitySweepStartRequest,
background_tasks: BackgroundTasks,
db: Annotated[Session, Depends(get_db)],
) -> CitySweepStartResponse:
"""Запустить Yandex.Недвижимость ЕКБ city sweep в background. Returns run_id.
Один прогон: SERP (fetch_around_multi_room, rooms × price-range combos)
→ save_listings → address-enrich (detail <title> → полный адрес с домом).
Прокси уже внутри YandexRealtyScraper + address-enrich сессии.
Coop cancel: POST /scrape/yandex-city-sweep/{run_id}/cancel.
НЕ добавляется в scrape_schedules — ручной/контролируемый запуск
(schedule yandex_city_sweep DORMANT, включается оператором вручную).
"""
run_id = runs_mod.create_run(db, source="yandex_city_sweep", params=payload.model_dump())
async def _sweep_task() -> None:
sweep_db = SessionLocal()
try:
await run_yandex_city_sweep(
sweep_db,
run_id=run_id,
config=RealScraperConfig(),
matcher=RealMatcherAdapter(),
enrichment=RealEnrichmentJobs(),
proxy_provider=_kit_proxy_provider(),
radius_m=payload.radius_m,
pages_per_anchor=payload.pages_per_anchor,
request_delay_sec=payload.request_delay_sec,
enrich_address=payload.enrich_address,
)
except Exception:
logger.exception("yandex-sweep background task run_id=%d crashed", run_id)
finally:
sweep_db.close()
background_tasks.add_task(_sweep_task)
logger.info("yandex-sweep queued run_id=%d params=%s", run_id, payload.model_dump())
return CitySweepStartResponse(
run_id=run_id,
status="running",
pages_per_anchor=payload.pages_per_anchor,
detail_top_n=0, # у Yandex sweep нет detail-шага (address-enrich отдельно)
)
@router.get("/scrape/yandex-city-sweep/runs", response_model=list[ScrapeRunRow])
def list_yandex_city_sweep_runs(
db: Annotated[Session, Depends(get_db)],
limit: int = 10,
) -> list[ScrapeRunRow]:
"""Список последних N Yandex city sweep runs (для UI polling). Default limit=10."""
rows = runs_mod.list_recent(db, source="yandex_city_sweep", limit=limit)
return [
ScrapeRunRow(
run_id=r["run_id"],
source=r["source"],
status=r["status"],
params=r.get("params"),
counters=r.get("counters"),
error=r.get("error"),
started_at=r["started_at"].isoformat() if r.get("started_at") else None,
finished_at=r["finished_at"].isoformat() if r.get("finished_at") else None,
heartbeat_at=r["heartbeat_at"].isoformat() if r.get("heartbeat_at") else None,
)
for r in rows
]
@router.post("/scrape/yandex-city-sweep/{run_id}/cancel")
def cancel_yandex_city_sweep(
run_id: int,
db: Annotated[Session, Depends(get_db)],
) -> dict[str, object]:
"""Отменить running Yandex sweep. Cooperative: проверяется каждый anchor."""
cancelled = runs_mod.mark_cancelled(db, run_id)
return {"ok": True, "run_id": run_id, "cancelled": cancelled}
# ── Scraper settings: global + per-source delay management ───────────────────
class ScraperSettingPayload(BaseModel):
source: str = Field(..., min_length=1, max_length=64)
request_delay_sec: float = Field(..., ge=0.0, le=60.0)
description: str | None = None
@router.get("/scraper-settings", status_code=200)
def list_scraper_settings(
db: Annotated[Session, Depends(get_db)],
) -> dict:
"""Список всех настроек задержки (per-source + global)."""
rows = (
db.execute(
text("""
SELECT source, request_delay_sec, description, updated_at
FROM scraper_settings
ORDER BY source ASC
""")
)
.mappings()
.all()
)
return {"settings": [dict(r) for r in rows]}
@router.put("/scraper-settings/{source}", status_code=200)
def update_scraper_setting(
source: str,
payload: ScraperSettingPayload,
db: Annotated[Session, Depends(get_db)],
) -> dict:
"""Обновить задержку для source (или 'global'). Path source используется как ключ.
Эффект применяется немедленно — invalidate_cache() сбрасывает кеш для source,
следующий get_scraper_delay() перечитает из БД.
source='global' — нижняя планка для всех парсеров (0 = отключено).
"""
if payload.source != source:
raise HTTPException(
status_code=400,
detail=f"Path source={source!r} does not match body source={payload.source!r}",
)
row = (
db.execute(
text("""
INSERT INTO scraper_settings (source, request_delay_sec, description, updated_at)
VALUES (:s, CAST(:d AS numeric), :desc, NOW())
ON CONFLICT (source) DO UPDATE SET
request_delay_sec = EXCLUDED.request_delay_sec,
description = COALESCE(EXCLUDED.description, scraper_settings.description),
updated_at = NOW()
RETURNING source, request_delay_sec, description, updated_at
"""),
{
"s": source,
"d": payload.request_delay_sec,
"desc": payload.description,
},
)
.mappings()
.one()
)
db.commit()
invalidate_cache(source)
logger.info("scraper-settings: updated source=%s delay=%.1f", source, payload.request_delay_sec)
return dict(row)
# ── In-app scheduler endpoints (Stage 4e) ────────────────────────────────────
@router.get("/scrape/schedules", response_model=list[ScheduleConfig])
def list_schedules(db: Annotated[Session, Depends(get_db)]) -> list[ScheduleConfig]:
"""Список всех schedules. Сейчас single row 'avito_city_sweep', extensible."""
rows = (
db.execute(
text(
"""
SELECT id, source, enabled, window_start_hour, window_end_hour,
default_params, last_run_id, last_run_at, next_run_at,
created_at, updated_at
FROM scrape_schedules
ORDER BY source
"""
),
)
.mappings()
.all()
)
return [
ScheduleConfig(
id=r["id"],
source=r["source"],
enabled=r["enabled"],
window_start_hour=r["window_start_hour"],
window_end_hour=r["window_end_hour"],
default_params=r["default_params"] or {},
last_run_id=r["last_run_id"],
last_run_at=r["last_run_at"].isoformat() if r["last_run_at"] else None,
next_run_at=r["next_run_at"].isoformat() if r["next_run_at"] else None,
updated_at=r["updated_at"].isoformat() if r["updated_at"] else None,
)
for r in rows
]
@router.put("/scrape/schedules/{source}", response_model=ScheduleConfig)
def update_schedule(
source: str,
payload: ScheduleConfigUpdate,
db: Annotated[Session, Depends(get_db)],
) -> ScheduleConfig:
"""UPDATE existing schedule (create если не существует, через INSERT ON CONFLICT)."""
from app.services.scheduler import compute_next_run_at
# #2674: такт берётся из default_params — ровно как его читает планировщик
# (_claim_run/_defer_next_run_at). Без него compute_next_run_at падал на default=1 и
# ЛЮБОЕ сохранение сбивало источник на «завтра»: недельный avito_full_load после
# правки окна побежал бы через сутки. На суточных источниках баг был невидим —
# для них «завтра» и есть правильный ответ.
# None-safe так же, как в scheduler: `"interval_days": null` в jsonb → 1, не TypeError.
_interval_days = payload.default_params.get("interval_days")
interval_days = max(1, int(_interval_days)) if _interval_days is not None else 1
# Явно заданный оператором момент уважается как есть (в т.ч. в прошлом — «запустить
# сейчас»). Иначе считаем от такта.
next_at = payload.next_run_at or compute_next_run_at(
payload.window_start_hour,
payload.window_end_hour,
interval_days=interval_days,
)
row = (
db.execute(
text(
"""
INSERT INTO scrape_schedules
(source, enabled, window_start_hour, window_end_hour, default_params, next_run_at)
VALUES (:source, :enabled, :ws, :we, CAST(:params AS jsonb), :next_at)
ON CONFLICT (source) DO UPDATE SET
enabled = EXCLUDED.enabled,
window_start_hour = EXCLUDED.window_start_hour,
window_end_hour = EXCLUDED.window_end_hour,
default_params = EXCLUDED.default_params,
-- #2674: не двигаем уже назначенный запуск, если двигать не за чем.
-- Раньше next_run_at перезаписывался ВСЕГДА, поэтому правка соседнего
-- поля (enabled, request_delay_sec в params) заново разыгрывала момент
-- внутри окна и сдвигала прогон. Сохраняем существующий только когда он
-- ещё в будущем И ни окно, ни такт не менялись — тогда пересчёт дал бы
-- то же самое окно, только с другим random-смещением.
next_run_at = CASE
WHEN CAST(:explicit AS boolean) THEN EXCLUDED.next_run_at
WHEN scrape_schedules.next_run_at > NOW()
AND scrape_schedules.window_start_hour = EXCLUDED.window_start_hour
AND scrape_schedules.window_end_hour = EXCLUDED.window_end_hour
AND COALESCE(scrape_schedules.default_params ->> 'interval_days', '1')
= COALESCE(EXCLUDED.default_params ->> 'interval_days', '1')
THEN scrape_schedules.next_run_at
ELSE EXCLUDED.next_run_at
END,
updated_at = NOW()
RETURNING id, source, enabled, window_start_hour, window_end_hour,
default_params, last_run_id, last_run_at, next_run_at, updated_at
"""
),
{
"source": source,
"enabled": payload.enabled,
"ws": payload.window_start_hour,
"we": payload.window_end_hour,
"params": json.dumps(payload.default_params, ensure_ascii=False),
"next_at": next_at,
"explicit": payload.next_run_at is not None,
},
)
.mappings()
.fetchone()
)
db.commit()
assert row is not None, f"scrape_schedules upsert returned no row for source={source!r}"
return ScheduleConfig(
id=row["id"],
source=row["source"],
enabled=row["enabled"],
window_start_hour=row["window_start_hour"],
window_end_hour=row["window_end_hour"],
default_params=row["default_params"] or {},
last_run_id=row["last_run_id"],
last_run_at=row["last_run_at"].isoformat() if row["last_run_at"] else None,
next_run_at=row["next_run_at"].isoformat() if row["next_run_at"] else None,
updated_at=row["updated_at"].isoformat() if row["updated_at"] else None,
)
# -- Yandex ad-hoc scrape triggers --------------------------------------------
class YandexDetailTriggerResp(BaseModel):
ok: bool
offer_url: str
offer_id: str | None
price_rub: int | None
title: str | None
photo_count: int
@router.post("/scrape/yandex-detail", response_model=YandexDetailTriggerResp)
async def scrape_yandex_detail(
offer_url: str,
) -> YandexDetailTriggerResp:
"""Ad-hoc parse one Yandex offer detail page.
Returns a snapshot of the extracted DetailEnrichment (no DB write - read-only
debug; main pipeline writes via estimator on /estimate flow).
"""
_assert_allowed_url(offer_url)
async with YandexDetailScraper(delay_provider=get_scraper_delay) as scraper:
result = await scraper.fetch_detail(offer_url)
if result is None:
raise HTTPException(404, f"Could not parse Yandex detail page: {offer_url}")
return YandexDetailTriggerResp(
ok=True,
offer_url=offer_url,
offer_id=result.offer_id,
price_rub=result.price_rub,
title=result.title,
photo_count=len(result.photo_urls),
)
class YandexNewbuildingTriggerResp(BaseModel):
ok: bool
ext_id: str
ext_slug: str
name: str | None
lat: float | None
lon: float | None
rating: float | None
ratings_count: int | None
text_reviews_count: int | None
developer_name: str | None
@router.post("/scrape/yandex-newbuilding", response_model=YandexNewbuildingTriggerResp)
async def scrape_yandex_newbuilding(
slug: str,
id: str,
city: str = "ekaterinburg",
) -> YandexNewbuildingTriggerResp:
"""Ad-hoc parse a Yandex JK landing page.
URL: /{city}/kupit/novostrojka/<slug>-<id>/
"""
# fetch_jk использует внутренний BrowserFetcher, httpx-клиент BaseScraper не нужен.
# config= обязателен (#2322 fix) — иначе BrowserFetcher строится без endpoint=.
scraper = YandexNewbuildingScraper(config=RealScraperConfig())
result = await scraper.fetch_jk(jk_slug=slug, jk_id=id, city=city)
if result is None:
raise HTTPException(404, f"Could not parse Yandex JK: {slug}-{id} in {city}")
return YandexNewbuildingTriggerResp(
ok=True,
ext_id=result.ext_id,
ext_slug=result.ext_slug,
name=result.name,
lat=result.lat,
lon=result.lon,
rating=result.rating,
ratings_count=result.ratings_count,
text_reviews_count=result.text_reviews_count,
developer_name=result.developer_name,
)
class YandexNewbuildingSweepStartRequest(BaseModel):
limit: int = Field(default=5, ge=1, le=200)
request_delay_sec: float = Field(default=8.0, ge=3.0, le=30.0)
force: bool = False
city: str = Field(default="ekaterinburg", min_length=1, max_length=64)
class YandexNewbuildingSweepStartResponse(BaseModel):
run_id: int
status: str
limit: int
detail: str
@router.post("/scrape/yandex-newbuilding-sweep", response_model=YandexNewbuildingSweepStartResponse)
async def start_yandex_newbuilding_sweep(
payload: YandexNewbuildingSweepStartRequest,
background_tasks: BackgroundTasks,
db: Annotated[Session, Depends(get_db)],
) -> YandexNewbuildingSweepStartResponse:
"""Запустить Yandex newbuilding enrichment sweep в background (#974). Returns run_id.
Один прогон: SELECT pending yandex_realty_nb houses → resolve slug (если NULL) →
fetch_jk через BrowserFetcher → UPSERT market.yandex_jk_enrichment.
Shipped DORMANT (scheduler seed enabled=false). Этот endpoint — для ручного запуска.
Coop cancel через DELETE scrape_runs или cancel endpoint.
"""
from app.tasks.yandex_newbuilding_sweep import enrich_yandex_newbuilding_sweep
run_id = runs_mod.create_run(db, source="yandex_newbuilding_sweep", params=payload.model_dump())
async def _sweep_task() -> None:
sweep_db = SessionLocal()
try:
result = await enrich_yandex_newbuilding_sweep(
sweep_db,
limit=payload.limit,
force=payload.force,
request_delay_sec=payload.request_delay_sec,
city=payload.city,
)
runs_mod.mark_done(sweep_db, run_id, result.to_dict())
except Exception:
logger.exception("yandex-nb-sweep background task run_id=%d crashed", run_id)
try:
runs_mod.mark_failed(sweep_db, run_id, "crashed", {})
except Exception:
pass
finally:
sweep_db.close()
background_tasks.add_task(_sweep_task)
logger.info("yandex-nb-sweep queued run_id=%d params=%s", run_id, payload.model_dump())
return YandexNewbuildingSweepStartResponse(
run_id=run_id,
status="running",
limit=payload.limit,
detail=f"yandex_newbuilding_sweep started: limit={payload.limit} city={payload.city}",
)
@router.get("/scrape/yandex-newbuilding-sweep/runs", response_model=list[ScrapeRunRow])
def list_yandex_newbuilding_sweep_runs(
db: Annotated[Session, Depends(get_db)],
limit: int = 10,
) -> list[ScrapeRunRow]:
"""Список последних N yandex newbuilding sweep runs (для UI polling). Default limit=10."""
rows = runs_mod.list_recent(db, source="yandex_newbuilding_sweep", limit=limit)
return [
ScrapeRunRow(
run_id=r["run_id"],
source=r["source"],
status=r["status"],
params=r.get("params"),
counters=r.get("counters"),
error=r.get("error"),
started_at=r["started_at"].isoformat() if r.get("started_at") else None,
finished_at=r["finished_at"].isoformat() if r.get("finished_at") else None,
heartbeat_at=r["heartbeat_at"].isoformat() if r.get("heartbeat_at") else None,
)
for r in rows
]
class YandexValuationTriggerResp(BaseModel):
ok: bool
address: str
year_built: int | None
total_floors: int | None
house_type: str | None
has_lift: bool | None
total_objects: int | None
history_items_count: int
@router.post("/scrape/yandex-valuation", response_model=YandexValuationTriggerResp)
async def scrape_yandex_valuation(
address: str,
offer_category: str = "APARTMENT",
offer_type: str = "SELL",
page: int = 1,
) -> YandexValuationTriggerResp:
"""Ad-hoc fetch Yandex Valuation house-history for an address (debug, no cache write).
Anonymous GET - works without auth.
"""
async with YandexValuationScraper(
RealScraperConfig(), delay_provider=get_scraper_delay, proxy_provider=_kit_proxy_provider()
) as scraper:
result = await scraper.fetch_house_history(
address=address,
offer_category=offer_category,
offer_type=offer_type,
page=page,
)
if result is None:
raise HTTPException(404, f"Could not parse Yandex Valuation: {address}")
return YandexValuationTriggerResp(
ok=True,
address=address,
year_built=result.house.year_built,
total_floors=result.house.total_floors,
house_type=result.house.house_type,
has_lift=result.house.has_lift,
total_objects=result.house.total_objects,
history_items_count=len(result.history_items),
)
class CianDetailTriggerResp(BaseModel):
ok: bool
offer_url: str
cian_id: int | None
price_changes_count: int
saved: bool
listing_id: int | None
@router.post("/scrape/cian-detail", response_model=CianDetailTriggerResp)
async def scrape_cian_detail(
offer_url: str,
listing_id: int | None = None,
db: Session = Depends(get_db), # noqa: B008
) -> CianDetailTriggerResp:
"""Ad-hoc parse one Cian offer detail page.
If `listing_id` provided → save enrichment (offer_price_history + listings updates).
Without it → debug-only (no DB write).
"""
_assert_allowed_url(offer_url)
from scraper_kit.providers.cian.detail import fetch_detail, save_detail_enrichment
enrichment = await fetch_detail(offer_url, config=RealScraperConfig())
if enrichment is None:
raise HTTPException(404, f"Could not parse Cian detail page: {offer_url}")
saved = False
if listing_id is not None:
save_detail_enrichment(db, listing_id, enrichment, matcher=RealMatcherAdapter())
saved = True
return CianDetailTriggerResp(
ok=True,
offer_url=offer_url,
cian_id=enrichment.cian_id,
price_changes_count=len(enrichment.price_changes),
saved=saved,
listing_id=listing_id,
)
class CianNewbuildingTriggerResp(BaseModel):
ok: bool
zhk_url: str
cian_internal_house_id: int | None
name: str | None
price_dynamics_count: int
has_management_company: bool
saved: bool
house_id: int | None
@router.post("/scrape/cian-newbuilding", response_model=CianNewbuildingTriggerResp)
async def scrape_cian_newbuilding(
zhk_url: str,
house_id: int | None = None,
db: Session = Depends(get_db), # noqa: B008
) -> CianNewbuildingTriggerResp:
"""Ad-hoc parse a Cian ЖК (newbuilding) catalog page.
If `house_id` provided → persist (houses + houses_price_dynamics + management_companies).
Without it → debug-only (no DB write).
"""
_assert_allowed_url(zhk_url)
from scraper_kit.providers.cian.newbuilding import (
fetch_newbuilding,
save_newbuilding_enrichment,
)
enrichment = await fetch_newbuilding(
zhk_url, config=RealScraperConfig(), proxy_provider=_kit_proxy_provider()
)
if enrichment is None:
raise HTTPException(404, f"Could not parse Cian newbuilding page: {zhk_url}")
saved = False
if house_id is not None:
# save_newbuilding_enrichment — sync (def, returns None); await на sync-функции
# раньше поднимал TypeError на любом вызове с house_id.
save_newbuilding_enrichment(db, house_id, enrichment)
saved = True
return CianNewbuildingTriggerResp(
ok=True,
zhk_url=zhk_url,
cian_internal_house_id=enrichment.cian_internal_house_id,
name=enrichment.name,
price_dynamics_count=len(enrichment.realty_valuation_chart),
has_management_company=bool(enrichment.management_company),
saved=saved,
house_id=house_id,
)
class CianBackfillResp(BaseModel):
ok: bool
listings_total: int
listings_processed: int
listings_succeeded: int
listings_failed_fetch: int
listings_failed_save: int
price_changes_attempted: int
houses_total: int
houses_processed: int
houses_succeeded: int
houses_failed_fetch: int
houses_failed_save: int
valuations_total: int
valuations_processed: int
valuations_succeeded: int
valuations_failed: int
duration_sec: float
@router.post("/scrape/cian-backfill-history", response_model=CianBackfillResp)
async def scrape_cian_backfill_history(
batch_size: int = Query(50, ge=1, le=200),
do_listings: bool = True,
do_houses: bool = True,
do_valuations: bool = False,
dry_run: bool = False,
db: Session = Depends(get_db), # noqa: B008
) -> CianBackfillResp:
"""Batch backfill для Cian historical data.
Iterates Cian listings with missing offer_price_history, fetches Cian detail
pages, persists enrichment (price changes, views, agent, ceiling_height, etc.).
Houses block (houses_price_dynamics): fetches newbuilding pages for houses with
cian_zhk_url set and no existing price dynamics rows. Requires migration
071_houses_cian_zhk_url.sql applied and cian_zhk_url populated (via SERP scrape
or scripts/backfill_cian_zhk_url.py).
Rate-limit: scraper_settings 'cian' delay (default 5s) between requests.
Idempotent: skips listings/houses that already have history rows.
"""
from app.tasks.cian_history_backfill import CianBackfillResult, backfill_cian_history
result: CianBackfillResult = await backfill_cian_history(
db,
batch_size=batch_size,
do_listings=do_listings,
do_houses=do_houses,
do_valuations=do_valuations,
dry_run=dry_run,
)
return CianBackfillResp(ok=True, **result.__dict__)
# ── T8a: Cian price-history backfill ─────────────────────────────────────────
class CianPriceHistoryRequest(BaseModel):
batch_size: int = Field(default=50, ge=1, le=500)
listing_id: int | None = Field(
default=None,
description="Если указан — обработать один конкретный листинг.",
)
class BackfillCountersResp(BaseModel):
ok: bool
checked: int
saved: int
skipped: int
errors: int
duration_sec: float
@router.post("/scrape/cian-price-history", response_model=BackfillCountersResp)
async def scrape_cian_price_history(
payload: CianPriceHistoryRequest,
db: Session = Depends(get_db), # noqa: B008
) -> BackfillCountersResp:
"""T8a: Backfill offer_price_history для Cian-листингов.
Fetches Cian detail pages (curl_cffi, без Playwright) для листингов без строк
в offer_price_history, парсит priceChanges из _cianConfig defaultState.
batch_size: сколько листингов обработать за запуск (default 50).
listing_id: обработать один конкретный листинг (debug / точечная перезаливка).
Rate-limit: задержка из scraper_settings['cian'] (default 5s) между страницами.
Idempotent: повторный запуск безопасен (ON CONFLICT DO NOTHING).
"""
from app.services.cian_price_history import (
CianPriceHistoryResult,
backfill_cian_price_history,
)
result: CianPriceHistoryResult = await backfill_cian_price_history(
db,
batch_size=payload.batch_size,
listing_id=payload.listing_id,
)
return BackfillCountersResp(
ok=True,
checked=result.checked,
saved=result.saved,
skipped=result.skipped,
errors=result.errors,
duration_sec=round(result.duration_sec, 2),
)
# ── T10: Yandex address backfill ─────────────────────────────────────────────
class YandexAddressBackfillRequest(BaseModel):
limit: int = Field(default=200, ge=1, le=2000)
request_delay_sec: float = Field(default=3.0, ge=1.0, le=15.0)
@router.post("/scrape/yandex-address-backfill", response_model=BackfillCountersResp)
async def scrape_yandex_address_backfill(
payload: YandexAddressBackfillRequest,
background_tasks: BackgroundTasks,
db: Session = Depends(get_db), # noqa: B008
) -> BackfillCountersResp:
"""T10: Backfill listings.address для Yandex-листингов без номера дома.
Yandex SERP-карточки содержат только улицу («улица Горького»). Этот endpoint
фетчит detail-страницы через curl_cffi (chrome120), извлекает полный адрес
из <title> («Екатеринбург, улица Горького, 36 — id …»), обновляет
listings.address и сбрасывает geocode_tried_at для перегеокодирования.
limit: сколько листингов обработать за запуск (default 200).
request_delay_sec: пауза между запросами (default 3.0s).
Запускает задачу в foreground (не в BackgroundTasks), т.к. batch небольшой.
Для > 500 листингов рекомендуется вызывать несколько раз.
"""
from app.services.yandex_address_backfill import (
YandexAddressBackfillResult,
backfill_yandex_addresses,
)
result: YandexAddressBackfillResult = await backfill_yandex_addresses(
db,
limit=payload.limit,
request_delay_sec=payload.request_delay_sec,
)
return BackfillCountersResp(
ok=True,
checked=result.checked,
saved=result.saved,
skipped=result.skipped,
errors=result.errors,
duration_sec=round(result.duration_sec, 2),
)
# ── T7: House IMV bulk backfill ───────────────────────────────────────────────
class HouseIMVBackfillRequest(BaseModel):
batch_size: int = Field(default=50, ge=1, le=500)
request_delay_sec: float = Field(
default=5.0,
ge=1.0,
le=30.0,
description="Пауза между Avito IMV вызовами (default 5s — anti-captcha).",
)
only_status: str = Field(
default="pending",
description="Обрабатывать дома с этим imv_status. 'transient_error' — retry.",
)
house_id: int | None = Field(
default=None,
description="Обработать один конкретный дом (debug).",
)
class HouseIMVBackfillResp(BaseModel):
ok: bool
checked: int
saved: int
skipped: int
errors: int
duration_sec: float
status_counts: dict[str, int]
@router.post("/scrape/house-imv-backfill", response_model=HouseIMVBackfillResp)
async def scrape_house_imv_backfill(
payload: HouseIMVBackfillRequest,
background_tasks: BackgroundTasks,
db: Session = Depends(get_db), # noqa: B008
) -> HouseIMVBackfillResp:
"""T7: Avito IMV bulk backfill на уровне домов (house_imv_evaluations).
Итерирует дома с imv_status='pending' (или 'transient_error'), вычисляет
медианные параметры квартиры из связанных листингов, вызывает Avito IMV API,
сохраняет результаты в house_imv_evaluations + house_placement_history +
house_suggestions.
Resumable: дома помечаются imv_status='ok'/'not_found'/'error' после обработки.
Idempotent: повторный запуск не дублирует данные (UPSERT on house_id).
batch_size: сколько домов обработать за запуск (default 50).
request_delay_sec: пауза между IMV-вызовами (default 5s). ВАЖНО: Avito IMV
реагирует на частые запросы с datacenter-IP. Не снижать < 3s.
only_status: по умолчанию 'pending'. Для retry failed — 'transient_error'.
house_id: обработать один дом (debug).
Примечание по прокси: Avito IMV использует собственную curl_cffi-сессию.
settings.scraper_proxy_url прокидывается в неё автоматически через avito_imv.py
(evaluate_via_imv создаёт сессию без явного proxy-параметра — сессия использует
env HTTP_PROXY/HTTPS_PROXY если задан. При необходимости — добавить явный
proxy-параметр в evaluate_via_imv / backfill_house_imv на уровне сервиса).
"""
from app.services.house_imv_backfill import (
HouseIMVBackfillResult,
backfill_house_imv,
)
result: HouseIMVBackfillResult = await backfill_house_imv(
db,
batch_size=payload.batch_size,
request_delay_sec=payload.request_delay_sec,
only_status=payload.only_status,
house_id=payload.house_id,
)
return HouseIMVBackfillResp(
ok=True,
checked=result.checked,
saved=result.saved,
skipped=result.skipped,
errors=result.errors,
duration_sec=round(result.duration_sec, 2),
status_counts=result.status_counts,
)
# ── Единая scrapers-страница: unified runs + health (epic) ───────────────────
# rotate-ip (changeip mobileproxy) удалён #2616 шаг 3 — мёртвая подписка (#2613).
class UnifiedScrapeRunRow(BaseModel):
"""Строка scrape_runs для unified-таблицы (все source'ы в одной выдаче).
#2674: поля run_type больше нет. Вид прогона в БД всегда был дефолтом
'city_sweep' (3244 из 3244 строк, ни одно место кода его не задавало), и
таблица подписывала им прогоны, которые никаким sweep не были —
proxy_healthcheck, deactivate_stale_*, sber_index_pull. Что именно бежало,
называет `source`.
"""
run_id: int
source: str
status: str
# #2674: чинить фильтр без этого флага было бы регрессом. Пока таблица была
# пуста на всех вкладках, кнопка отмены не рендерилась ни разу; теперь оператор
# видит все 53 источника — и без флага мог бы «отменить» задачу, которая отмену
# не опрашивает (см. scrape_runs.honors_cancel): статус соврал бы, а
# has_running_run перестал бы держать single-run guard.
cancellable: bool = False
# #2686: диагноз для status='banned' — 'platform' (площадка заблокировала),
# 'infra' (не отдал наш браузерный сайдкар) или 'unknown' (#2764 — причина не
# установлена; раньше такие прогоны молча получали 'platform'). Без него
# оператор видит только «забанен» и делает вывод «площадка нас палит» на 80%
# наших же отказов.
ban_kind: str | None = None
params: dict | None = None
counters: dict | None = None
total_seen: int | None = None
new_count: int | None = None
started_at: str | None = None
finished_at: str | None = None
heartbeat_at: str | None = None
error_text: str | None = None
class UnifiedScrapeRunsResponse(BaseModel):
total: int
rows: list[UnifiedScrapeRunRow]
class ScrapeRunSourcesResponse(BaseModel):
"""Список source'ов для фильтра истории прогонов — из данных, не из литерала."""
sources: list[str]
class BrowserHealth(BaseModel):
reachable: bool
browsers: dict[str, bool] = Field(default_factory=dict)
class ProviderHealth(BaseModel):
source: str
proxy_host: str | None = None
proxy_port: int | None = None
rotate_supported: bool = False
current_ip: str | None = None
class ScraperHealthResponse(BaseModel):
fetch_mode: str
browser: BrowserHealth
providers: list[ProviderHealth]
_ROTATABLE_SOURCES = ("avito", "cian", "yandex")
def _provider_proxy_url(source: str) -> str | None:
"""Effective proxy URL для source (учитывает property-fallback в settings).
#2616 шаг 2: avito/cian/yandex все три сходятся на settings.scraper_proxy_url
(per-provider AVITO_PROXY_URL/CIAN_PROXY_URL/YANDEX_PROXY_URL сняты — мёртвая
mobileproxy-подписка, #2613).
"""
return {
"avito": settings.scraper_proxy_url,
"cian": settings.cian_proxy_url,
"yandex": settings.yandex_proxy_url,
}.get(source)
def _parse_proxy_host_port(proxy_url: str | None) -> tuple[str | None, int | None]:
"""Распарсить host/port из proxy URL (схема http(s)://user:pass@host:port)."""
if not proxy_url:
return None, None
try:
parsed = urlparse(proxy_url)
return parsed.hostname, parsed.port
except Exception:
logger.warning("scraper/health: cannot parse proxy url", exc_info=True)
return None, None
@router.get("/scrape/runs", response_model=UnifiedScrapeRunsResponse)
def list_scrape_runs_unified(
db: Annotated[Session, Depends(get_db)],
source: Annotated[str | None, Query()] = None,
status: Annotated[
# 'skipped' (#2658) — пропущенное расписание; без него оператор не может
# спросить «что сейчас пропускается» (фильтр отдавал 422 на единственной
# поверхности, построенной ровно для этого вопроса).
Literal["done", "running", "banned", "zombie", "failed", "cancelled", "skipped"] | None,
Query(),
] = None,
limit: Annotated[int, Query(ge=1, le=200)] = 50,
offset: Annotated[int, Query(ge=0)] = 0,
) -> UnifiedScrapeRunsResponse:
"""Unified-список scrape_runs по всем source'ам (для единой scrapers-страницы).
Замена 3 per-source /runs (avito/cian/yandex-city-sweep) единым эндпоинтом.
Старые per-source эндпоинты сохранены для обратной совместимости.
Query:
source — опц. фильтр по source (avito_city_sweep / cian_city_sweep / ...).
status — опц. фильтр (done/running/banned/zombie/failed/cancelled/skipped).
limit — default 50, max 200.
offset — default 0.
"""
total, rows = runs_mod.list_all(db, source=source, status=status, limit=limit, offset=offset)
def _iso(v: Any) -> str | None:
return v.isoformat() if v is not None else None
return UnifiedScrapeRunsResponse(
total=total,
rows=[
UnifiedScrapeRunRow(
run_id=r["run_id"],
source=r["source"],
status=r["status"],
cancellable=runs_mod.honors_cancel(str(r["source"])),
ban_kind=r.get("ban_kind"),
params=r.get("params"),
counters=r.get("counters"),
total_seen=r.get("total_seen"),
new_count=r.get("new_count"),
started_at=_iso(r.get("started_at")),
finished_at=_iso(r.get("finished_at")),
heartbeat_at=_iso(r.get("heartbeat_at")),
error_text=r.get("error_text"),
)
for r in rows
],
)
@router.get("/scrape/runs/sources", response_model=ScrapeRunSourcesResponse)
def list_scrape_run_sources(
db: Annotated[Session, Depends(get_db)],
) -> ScrapeRunSourcesResponse:
"""Источники для фильтра истории прогонов — ровно те, что есть в scrape_runs.
#2674: фильтр в UI был захардкожен тремя значениями (avito/cian/yandex), а в
таблице 53 разных source и НИ ОДНОЙ строки с таким точным значением — каждый
пункт фильтра давал пустую выдачу, и 76% прогонов (вся площадка Домклик в том
числе) были недоступны для вопроса «что там происходит». Список берётся из
данных: новый source появляется в фильтре сам, без правки кода.
"""
return ScrapeRunSourcesResponse(sources=runs_mod.distinct_sources(db))
async def _probe_browser_health() -> BrowserHealth:
"""GET tradein-browser /health (timeout 5с). reachable=False при ошибке."""
url = f"{settings.browser_http_endpoint.rstrip('/')}/health"
try:
async with httpx.AsyncClient(timeout=5.0) as client:
resp = await client.get(url)
resp.raise_for_status()
data = resp.json()
browsers = data.get("browsers") or {}
if not isinstance(browsers, dict):
browsers = {}
return BrowserHealth(reachable=True, browsers={k: bool(v) for k, v in browsers.items()})
except Exception:
logger.warning("scraper/health: browser /health unreachable", exc_info=True)
return BrowserHealth(reachable=False, browsers={})
async def _probe_current_ip(proxy_url: str | None) -> str | None:
"""Best-effort: текущий exit-IP через прокси (ipify, timeout 8с). None при ошибке."""
if not proxy_url:
return None
try:
async with httpx.AsyncClient(proxy=proxy_url, timeout=8.0) as client:
resp = await client.get("https://api.ipify.org", params={"format": "json"})
resp.raise_for_status()
ip = resp.json().get("ip")
return str(ip) if ip else None
except Exception:
logger.warning("scraper/health: ipify probe failed", exc_info=True)
return None
@router.get("/scraper/health", response_model=ScraperHealthResponse)
async def scraper_health() -> ScraperHealthResponse:
"""Сводный health для единой scrapers-страницы: fetch_mode + browser + провайдеры.
- fetch_mode: settings.scraper_fetch_mode (curl_cffi / browser).
- browser: GET tradein-browser /health (reachable + per-browser ready-флаги).
- providers: для avito/cian/yandex — proxy host/port, rotate_supported
(#2616 шаг 2: всегда False — changeip mobileproxy-ротация снята, мёртвый
аккаунт #2613; живая ASocks-ротация — POST /admin/proxies/{id}/rotate, #2611,
не per-provider-source), best-effort current_ip (пробинг через прокси на ipify).
Все пробинги параллельны (asyncio.gather) и time-boxed — суммарно ≤10с.
"""
proxy_urls = {s: _provider_proxy_url(s) for s in _ROTATABLE_SOURCES}
browser, *ips = await asyncio.gather(
_probe_browser_health(),
*[_probe_current_ip(proxy_urls[s]) for s in _ROTATABLE_SOURCES],
)
ip_by_source = dict(zip(_ROTATABLE_SOURCES, ips, strict=True))
providers: list[ProviderHealth] = []
for source in _ROTATABLE_SOURCES:
host, port = _parse_proxy_host_port(proxy_urls[source])
providers.append(
ProviderHealth(
source=source,
proxy_host=host,
proxy_port=port,
rotate_supported=False,
current_ip=ip_by_source[source],
)
)
return ScraperHealthResponse(
fetch_mode=settings.scraper_fetch_mode,
browser=browser,
providers=providers,
)
# ── Pacing live-регулятор (GET/PUT /scraper/pacing) ──────────────────────────
class PacingProvider(BaseModel):
source: str
interval_s: float
env_default_s: float
class PacingResponse(BaseModel):
providers: list[PacingProvider]
class PacingUpdateRequest(BaseModel):
interval_s: float = Field(ge=0, le=120)
async def _proxy_browser_pacing_get() -> dict:
"""GET tradein-browser /pacing (timeout 5с). Пробрасывает ответ или кидает HTTPException."""
url = f"{settings.browser_http_endpoint.rstrip('/')}/pacing"
try:
async with httpx.AsyncClient(timeout=5.0) as client:
resp = await client.get(url)
resp.raise_for_status()
return resp.json() # type: ignore[no-any-return]
except httpx.HTTPStatusError as exc:
raise HTTPException(
status_code=502,
detail=f"browser /pacing returned {exc.response.status_code}",
) from exc
except Exception as exc:
raise HTTPException(
status_code=503,
detail=f"browser unreachable: {type(exc).__name__}",
) from exc
@router.get("/scraper/pacing", response_model=PacingResponse)
async def get_scraper_pacing() -> PacingResponse:
"""Текущие live-интервалы пейсинга из tradein-browser (in-memory, per-provider).
interval_s — живой интервал (может быть изменён PUT).
env_default_s — значение из env (к чему сбросится на рестарте).
502/503 при недоступности browser-сервиса.
"""
data = await _proxy_browser_pacing_get()
providers = [
PacingProvider(
source=p["source"],
interval_s=float(p["interval_s"]),
env_default_s=float(p["env_default_s"]),
)
for p in data.get("providers", [])
]
return PacingResponse(providers=providers)
@router.put("/scraper/pacing/{source}", response_model=dict)
async def update_scraper_pacing(
source: Literal["avito", "cian", "yandex", "generic"],
payload: PacingUpdateRequest,
) -> dict:
"""Обновить live-интервал пейсинга для провайдера в tradein-browser (in-memory).
Значение сбрасывается к env-дефолту при рестарте контейнера — by design.
interval_s: 0..120 (сек). Возвращает {ok, source, interval_s}.
502/503 при недоступности browser-сервиса.
"""
url = f"{settings.browser_http_endpoint.rstrip('/')}/pacing"
try:
async with httpx.AsyncClient(timeout=5.0) as client:
resp = await client.put(
url,
json={"source": source, "interval_s": payload.interval_s},
)
resp.raise_for_status()
return resp.json() # type: ignore[no-any-return]
except httpx.HTTPStatusError as exc:
raise HTTPException(
status_code=502,
detail=f"browser /pacing returned {exc.response.status_code}",
) from exc
except Exception as exc:
raise HTTPException(
status_code=503,
detail=f"browser unreachable: {type(exc).__name__}",
) from exc
# ── Data-quality coverage (GET /scraper/data-quality) ────────────────────────
class SourceCoverage(BaseModel):
source: str
active_count: int
# #2660: «активно» ≠ «живо». is_active снимается только деактиватором протухших,
# а он покрывает не все источники — на проде (2026-08-05) cian показывал 18 530
# активных при 12 683 не виденных 14+ дней. Из-за этого #2574 месяц читалась как
# «всё собирается». Не прячем протухшее из счётчика, а отдаём ВТОРЫМ числом
# рядом — тогда «активно» перестаёт читаться как «живо».
stale_count: int
fields: dict[str, float] # field_name -> fill% (0..100, round 1)
class HousesCoverage(BaseModel):
total: int
validated_pct: float # avito_validated_at IS NOT NULL %
rating_pct: float # rating_score IS NOT NULL %
house_type_pct: float # house_type IS NOT NULL %
reviews_count: int
class DataQualityResponse(BaseModel):
sources: list[SourceCoverage]
houses: HousesCoverage
# Порог «не виделись N дней» для stale_count — отдаём в ответе, чтобы UI
# подписывал число, а не хардкодил порог у себя вторым определением.
stale_days: int
# Поля listings для fill%-аудита. Каждый кортеж: (имя_поля, SQL-выражение IS NOT NULL).
# Для photo_urls отдельная логика (jsonb не NULL и не '[]').
_DQ_LISTING_FIELDS: list[tuple[str, str]] = [
("description", "description IS NOT NULL AND description <> ''"),
("photo_urls", "photo_urls IS NOT NULL AND photo_urls <> '[]'::jsonb"),
("address", "address IS NOT NULL AND address <> ''"),
("lat", "lat IS NOT NULL"),
("lon", "lon IS NOT NULL"),
("kitchen_area_m2", "kitchen_area_m2 IS NOT NULL"),
("living_area_m2", "living_area_m2 IS NOT NULL"),
# #2699: одна колонка вместо двух. ceiling_height (019) DEPRECATED — писатели
# переведены на ceiling_height_m, исторические значения перенесены (мигр. 238).
("ceiling_height_m", "ceiling_height_m IS NOT NULL"),
("metro_stations", "metro_stations IS NOT NULL AND metro_stations <> '[]'::jsonb"),
]
@router.get("/scraper/data-quality", response_model=DataQualityResponse)
def get_data_quality(
db: Annotated[Session, Depends(get_db)],
) -> DataQualityResponse:
"""Fill%-аудит detail-полей listings по source + дом-статистика.
Один проход per source через COUNT(*)...FILTER — не N запросов.
Поля listings: description, photo_urls, address, lat/lon, kitchen_area_m2,
living_area_m2, ceiling_height_m (все источники, #2699), metro_stations.
houses: total, avito_validated_at%, rating_score%, house_type%.
house_reviews: общий count.
#2660: рядом с active_count отдаётся stale_count — сколько из «активных» не
виделись LISTINGS_FRESH_DAYS дней (last_seen_at). Порог отдаётся в ответе
(stale_days), чтобы UI не заводил второе определение.
"""
# Строим single-pass SELECT для listings полей через FILTER-агрегаты.
# Структура: COUNT(*) FILTER (WHERE <expr>) / NULLIF(COUNT(*), 0) * 100
# Одним запросом получаем active_count + все fill-counts по каждому source.
filter_exprs = ", ".join(
f"COUNT(*) FILTER (WHERE {expr}) AS f_{name}" for name, expr in _DQ_LISTING_FIELDS
)
# last_seen_at, а не scraped_at: счётчик отвечает буквально на «сколько не
# виделись». На проде две колонки не расходятся (замер 2026-08-05: 0 активных
# строк с разницей ≥ суток), но семантика счётчика — про «видели», и колонка
# должна называть ровно её.
sql_listings = text(f"""
SELECT
source,
COUNT(*) AS active_count,
COUNT(*) FILTER (
WHERE last_seen_at <= NOW() - (:fresh_days || ' days')::interval
) AS stale_count,
{filter_exprs}
FROM listings
WHERE is_active = true
GROUP BY source
ORDER BY source
""")
rows = db.execute(sql_listings, {"fresh_days": LISTINGS_FRESH_DAYS}).mappings().all()
sources: list[SourceCoverage] = []
for row in rows:
cnt = int(row["active_count"]) or 1 # защита от деления на ноль
fields: dict[str, float] = {}
for name, _ in _DQ_LISTING_FIELDS:
raw = row[f"f_{name}"]
fill_pct = round(float(raw or 0) / cnt * 100, 1)
fields[name] = fill_pct
sources.append(
SourceCoverage(
source=row["source"],
active_count=int(row["active_count"]),
stale_count=int(row["stale_count"] or 0),
fields=fields,
)
)
# Houses: один агрегат
sql_houses = text("""
SELECT
COUNT(*) AS total,
COUNT(*) FILTER (WHERE avito_validated_at IS NOT NULL) AS validated_cnt,
COUNT(*) FILTER (WHERE rating_score IS NOT NULL) AS rating_cnt,
COUNT(*) FILTER (WHERE house_type IS NOT NULL) AS house_type_cnt
FROM houses
""")
h = db.execute(sql_houses).mappings().one()
total_houses = int(h["total"]) or 1 # guard против деления на ноль в %
reviews_count_row = db.execute(text("SELECT COUNT(*) AS cnt FROM house_reviews")).scalar()
reviews_count = int(reviews_count_row or 0)
houses = HousesCoverage(
total=int(h["total"]),
validated_pct=round(float(h["validated_cnt"] or 0) / total_houses * 100, 1),
rating_pct=round(float(h["rating_cnt"] or 0) / total_houses * 100, 1),
house_type_pct=round(float(h["house_type_cnt"] or 0) / total_houses * 100, 1),
reviews_count=reviews_count,
)
return DataQualityResponse(sources=sources, houses=houses, stale_days=LISTINGS_FRESH_DAYS)
# ── Proxy pool: хранилище + bulk-загрузка / список (#2161) ───────────────────
#
# АДДИТИВНО: ручки лишь наполняют/читают scrape_proxies. Подключение пула к
# боевому сбору (pick/lease/rotate вместо env) — отдельные шаги P3/P4.
_PROXY_AFFINITIES = frozenset({"avito", "cian", "yandex", "domclick", "generic", "any"})
_PROXY_KINDS = frozenset({"http", "socks5"})
def _mask_proxy_url(url: str | None) -> str | None:
"""Маскирует пароль в proxy-URL: socks5://user:pass@host → socks5://user:***@host.
URL без userinfo/пароля возвращается как есть. Не-парсящийся URL — как есть
(лучше вернуть исходное, чем упасть на форматировании admin-листинга).
"""
if not url:
return url
try:
parsed = urlparse(url)
except ValueError:
return url
if not parsed.password:
return url
userinfo = parsed.username or ""
host = parsed.hostname or ""
masked_netloc = f"{userinfo}:***@{host}"
if parsed.port:
masked_netloc += f":{parsed.port}"
return urlunparse(parsed._replace(netloc=masked_netloc))
class ProxyIn(BaseModel):
"""Одна запись прокси для bulk-загрузки."""
url: str = Field(..., min_length=3, max_length=500)
provider_affinity: str = Field(default="any")
kind: str = Field(default="http")
rotate_url: str | None = Field(default=None, max_length=500)
label: str | None = Field(default=None, max_length=200)
geo: str | None = Field(default=None, max_length=100)
operator: str | None = Field(default=None, max_length=100)
@field_validator("provider_affinity")
@classmethod
def _check_affinity(cls, v: str) -> str:
if v not in _PROXY_AFFINITIES:
raise ValueError(
f"provider_affinity must be one of {sorted(_PROXY_AFFINITIES)}, got {v!r}"
)
return v
@field_validator("kind")
@classmethod
def _check_kind(cls, v: str) -> str:
if v not in _PROXY_KINDS:
raise ValueError(f"kind must be one of {sorted(_PROXY_KINDS)}, got {v!r}")
return v
class ProxyBulkRequest(BaseModel):
proxies: list[ProxyIn] = Field(..., min_length=1, max_length=1000)
class ProxyBulkResponse(BaseModel):
inserted: int
updated: int
class ProxySourceBan(BaseModel):
"""Активный бан узла КОНКРЕТНОЙ площадкой (#2600 п.2, scrape_proxy_source_bans)."""
source: str
banned_until: str
ban_count: int
def _fetch_source_bans(db: Session, proxy_ids: list[int]) -> dict[int, list[ProxySourceBan]]:
"""Активные (banned_until > now()) баны по источникам для указанных узлов.
Без этого оператор видит `enabled=true` и не понимает, почему узел не выдаётся
конкретному источнику (#2600 п.2 — бан теперь по паре «узел × источник», а не
глобальное выключение). Истёкшие строки не показываем: они ни на что не влияют,
живут ещё SOURCE_BAN_PURGE_DAYS только как память об эскалации.
"""
if not proxy_ids:
return {}
rows = (
db.execute(
text(
"""
SELECT proxy_id, source, banned_until, ban_count
FROM scrape_proxy_source_bans
WHERE banned_until > now()
AND proxy_id = ANY(CAST(:ids AS bigint[]))
ORDER BY proxy_id, source
"""
),
{"ids": proxy_ids},
)
.mappings()
.all()
)
bans: dict[int, list[ProxySourceBan]] = {}
for r in rows:
bans.setdefault(int(r["proxy_id"]), []).append(
ProxySourceBan(
source=r["source"],
banned_until=r["banned_until"].isoformat(),
ban_count=r["ban_count"],
)
)
return bans
class ProxyRow(BaseModel):
id: int
label: str | None
url: str # маскированный
kind: str
provider_affinity: str
rotate_url: str | None # маскированный
enabled: bool
disabled_reason: str | None # #2610: NULL = не выключен вручную (авто-воскрешаем)
consecutive_fails: int
exit_ip: str | None
latency_ms: int | None
last_check_at: str | None
last_ok_at: str | None
leased_by: int | None
leased_at: str | None
geo: str | None
operator: str | None
expires_at: str | None
created_at: str | None
updated_at: str | None
# #2600 п.2: активные баны площадками. Пустой список = узел выдаётся всем источникам.
source_bans: list[ProxySourceBan] = Field(default_factory=list)
@router.post("/proxies/bulk", response_model=ProxyBulkResponse)
def bulk_upsert_proxies(
payload: ProxyBulkRequest,
db: Annotated[Session, Depends(get_db)],
) -> ProxyBulkResponse:
"""Bulk UPSERT прокси по url (#2161).
Тело: {"proxies": [{"url", "provider_affinity", "kind"?, "rotate_url"?,
"label"?, "geo"?, "operator"?}, ...]}.
Существующий url → DO UPDATE (affinity/kind/rotate_url + enabled,
label/geo/operator обновляются если переданы). Новый → INSERT (enabled=true,
disabled_reason=NULL — новый прокси не может быть "выключен вручную").
enabled на UPDATE-ветке НЕ безусловный (#2610): если у существующей строки
disabled_reason НЕ NULL (оператор снял узел с ротации вручную), bulk-upsert
(например повторный прогон загрузчика с тем же url) не должен тихо вернуть
его в строй — тот же класс бага, что чинили в mark_health. enabled=true
ставится, только если disabled_reason IS NULL; сам disabled_reason bulk
не трогает (эта ручка не умеет ни ставить, ни снимать ручной флаг — это
PATCH /proxies/{id}, см. patch_proxy).
Валидация provider_affinity/kind по whitelist на уровне Pydantic → 422.
Возвращает {inserted, updated}. Дубли по url ВНУТРИ одного запроса
схлопываются последним вхождением (ON CONFLICT в рамках одной транзакции).
"""
inserted = 0
updated = 0
for proxy in payload.proxies:
# RETURNING (xmax = 0) — системный трюк Postgres: xmax=0 у только что
# вставленной строки, !=0 у обновлённой существующей.
row = db.execute(
text(
"""
INSERT INTO scrape_proxies
(url, provider_affinity, kind, rotate_url, label, geo, operator,
enabled, updated_at)
VALUES
(:url, :aff, :kind, :rotate_url, :label, :geo, :operator, true, now())
ON CONFLICT (url) DO UPDATE SET
provider_affinity = EXCLUDED.provider_affinity,
kind = EXCLUDED.kind,
rotate_url = EXCLUDED.rotate_url,
label = COALESCE(EXCLUDED.label, scrape_proxies.label),
geo = COALESCE(EXCLUDED.geo, scrape_proxies.geo),
operator = COALESCE(EXCLUDED.operator, scrape_proxies.operator),
enabled = CASE
WHEN scrape_proxies.disabled_reason IS NULL THEN true
ELSE scrape_proxies.enabled
END,
updated_at = now()
RETURNING (xmax = 0) AS was_inserted
"""
),
{
"url": proxy.url,
"aff": proxy.provider_affinity,
"kind": proxy.kind,
"rotate_url": proxy.rotate_url,
"label": proxy.label,
"geo": proxy.geo,
"operator": proxy.operator,
},
).scalar()
if row:
inserted += 1
else:
updated += 1
db.commit()
logger.info("proxies/bulk: inserted=%d updated=%d", inserted, updated)
return ProxyBulkResponse(inserted=inserted, updated=updated)
@router.get("/proxies", response_model=list[ProxyRow])
def list_proxies(
db: Annotated[Session, Depends(get_db)],
provider: str | None = Query(default=None, description="Фильтр по provider_affinity"),
enabled: bool | None = Query(default=None, description="Фильтр по enabled"),
) -> list[ProxyRow]:
"""Список прокси со статусами. Пароли в url/rotate_url маскируются.
Фильтры: provider (=provider_affinity), enabled. Без фильтров — все.
source_bans — активные баны узла площадками (#2600 п.2): узел может быть
enabled=true и при этом не выдаваться конкретному источнику.
"""
clauses: list[str] = []
params: dict[str, Any] = {}
if provider is not None:
clauses.append("provider_affinity = :provider")
params["provider"] = provider
if enabled is not None:
clauses.append("enabled = :enabled")
params["enabled"] = enabled
where = ("WHERE " + " AND ".join(clauses)) if clauses else ""
rows = (
db.execute(
text(
f"""
SELECT id, label, url, kind, provider_affinity, rotate_url, enabled,
disabled_reason, consecutive_fails, exit_ip, latency_ms,
last_check_at, last_ok_at, leased_by, leased_at, geo, operator,
expires_at, created_at, updated_at
FROM scrape_proxies
{where}
ORDER BY provider_affinity, id
"""
# where собран из фикс-строк, значения параметризованы (no injection)
),
params,
)
.mappings()
.all()
)
def _iso(v: Any) -> str | None:
return v.isoformat() if v is not None else None
bans = _fetch_source_bans(db, [int(r["id"]) for r in rows])
return [
ProxyRow(
id=r["id"],
label=r["label"],
url=_mask_proxy_url(r["url"]) or "",
kind=r["kind"],
provider_affinity=r["provider_affinity"],
rotate_url=_mask_proxy_url(r["rotate_url"]),
enabled=r["enabled"],
disabled_reason=r["disabled_reason"],
consecutive_fails=r["consecutive_fails"],
exit_ip=r["exit_ip"],
latency_ms=r["latency_ms"],
last_check_at=_iso(r["last_check_at"]),
last_ok_at=_iso(r["last_ok_at"]),
leased_by=r["leased_by"],
leased_at=_iso(r["leased_at"]),
geo=r["geo"],
operator=r["operator"],
expires_at=_iso(r["expires_at"]),
created_at=_iso(r["created_at"]),
updated_at=_iso(r["updated_at"]),
source_bans=bans.get(int(r["id"]), []),
)
for r in rows
]
class ProxyPatch(BaseModel):
enabled: bool
reason: str | None = Field(
default=None,
max_length=500,
description=(
"Причина ручного выключения (#2610). Используется только когда enabled=false; "
"при отсутствии подставляется дефолтный текст. Игнорируется при enabled=true — "
"включение всегда сбрасывает disabled_reason в NULL."
),
)
_DEFAULT_MANUAL_DISABLE_REASON = "manually disabled via admin API"
@router.patch("/proxies/{proxy_id}", response_model=ProxyRow)
def patch_proxy(
proxy_id: int,
payload: ProxyPatch,
db: Annotated[Session, Depends(get_db)],
) -> ProxyRow:
"""Enable/disable одного прокси по id. 404 если не найден.
#2610: разводит "ручное выключение оператором" от "авто-выключение пулом".
enabled=false → disabled_reason ставится (payload.reason либо дефолтный текст) —
mark_health(ok=True) больше не воскресит узел молча первой успешной ipify-пробой.
enabled=true → disabled_reason ОБЯЗАТЕЛЬНО сбрасывается в NULL — иначе узел,
однажды выключенный руками, никогда больше не участвовал бы в авто-восстановлении
(см. proxy_pool.mark_health).
"""
row = (
db.execute(
text(
"""
UPDATE scrape_proxies
SET enabled = :enabled,
disabled_reason = CASE
WHEN CAST(:enabled AS boolean) THEN NULL
ELSE COALESCE(CAST(:reason AS text), disabled_reason,
CAST(:default_reason AS text))
END,
updated_at = now()
WHERE id = :id
RETURNING id, label, url, kind, provider_affinity, rotate_url, enabled,
disabled_reason, consecutive_fails, exit_ip, latency_ms,
last_check_at, last_ok_at, leased_by, leased_at, geo, operator,
expires_at, created_at, updated_at
"""
),
{
"enabled": payload.enabled,
"reason": payload.reason,
"default_reason": _DEFAULT_MANUAL_DISABLE_REASON,
"id": proxy_id,
},
)
.mappings()
.fetchone()
)
if row is None:
raise HTTPException(status_code=404, detail=f"proxy id={proxy_id} not found")
db.commit()
if payload.enabled:
# Ручное включение = чистый лист, как и обнуление disabled_reason выше (#2610).
# Иначе узел вернулся бы enabled=true, но по-прежнему невыдаваемым источникам с
# активным баном — и оператор не имел бы способа снять ложный бан (#2600 п.2).
clear_source_bans(db, proxy_id, reason="manual enable via admin API")
if not payload.enabled:
logger.info(
"proxy_pool: proxy id=%d manually disabled via admin API (reason=%r) — "
"auto-revive suspended until re-enabled (#2610)",
proxy_id,
row["disabled_reason"],
)
def _iso(v: Any) -> str | None:
return v.isoformat() if v is not None else None
return ProxyRow(
id=row["id"],
label=row["label"],
url=_mask_proxy_url(row["url"]) or "",
kind=row["kind"],
provider_affinity=row["provider_affinity"],
rotate_url=_mask_proxy_url(row["rotate_url"]),
enabled=row["enabled"],
disabled_reason=row["disabled_reason"],
consecutive_fails=row["consecutive_fails"],
exit_ip=row["exit_ip"],
latency_ms=row["latency_ms"],
last_check_at=_iso(row["last_check_at"]),
last_ok_at=_iso(row["last_ok_at"]),
leased_by=row["leased_by"],
leased_at=_iso(row["leased_at"]),
geo=row["geo"],
operator=row["operator"],
expires_at=_iso(row["expires_at"]),
created_at=_iso(row["created_at"]),
updated_at=_iso(row["updated_at"]),
source_bans=_fetch_source_bans(db, [int(row["id"])]).get(int(row["id"]), []),
)
# ── Proxy pool: ручная ротация exit-IP по proxy_id (#2600 п.5) ───────────────
#
# Раньше отдельно от /scraper/{source}/rotate-ip (env-прокси mobileproxy,
# changeip-ссылка) — тот эндпоинт удалён вместе с мёртвой подпиской (#2616 шаг 3).
# Этот эндпоинт — единственная живая ручная ротация, по proxy_id из пула
# scrape_proxies (сейчас это ASocks-порты с суточным лимитом 3/сутки), см.
# app.services.proxy_rotation.rotate_proxy.
class ProxyRotateResponse(BaseModel):
ok: bool
reason: str | None = None
new_ip: str | None = None
rotations_remaining_today: int
@router.post("/proxies/{proxy_id}/rotate", response_model=ProxyRotateResponse)
async def rotate_pool_proxy(
proxy_id: int,
db: Annotated[Session, Depends(get_db)],
) -> ProxyRotateResponse:
"""Ручная ротация exit-IP одного прокси пула (#2600 п.5).
Делегирует в app.services.proxy_rotation.rotate_proxy — читает rotate_url
прокси из scrape_proxies, требует ASOCKS_API_TOKEN (settings.asocks_api_token),
проверяет суточный лимит (3/сутки, scrape_proxy_rotations) ДО обращения к API.
ok=False — ожидаемая бизнес-ситуация (нет rotate_url / нет токена / лимит /
провайдер отказал), НЕ HTTPException; reason ВСЕГДА нейтральный, без токена.
ПОКА без автотриггера по бану (issue #2600 п.2: сигнал бана до пула не
доходит — страница-заглушка отдаёт 200) — только этот ручной вызов.
"""
result = await proxy_rotation_svc.rotate_proxy(db, proxy_id)
return ProxyRotateResponse(
ok=result.ok,
reason=result.reason,
new_ip=result.new_ip,
rotations_remaining_today=result.rotations_remaining_today,
)