"""Admin endpoints — debug + ручной запуск парсеров. Пока без auth (Слой 7), потом закроем общим middleware. """ from __future__ import annotations import asyncio import logging from typing import Annotated, Literal from fastapi import APIRouter, Depends, Header, HTTPException from pydantic import BaseModel, Field from sqlalchemy import text from sqlalchemy.orm import Session from app.core.config import settings from app.core.db import get_db from app.services.geocoder import geocode from app.services.scrapers.avito import AvitoScraper from app.services.scrapers.base import save_listings from app.services.scrapers.cian import CianScraper from app.services.scrapers.n1 import N1Scraper from app.services.scrapers.yandex_realty import YandexRealtyScraper logger = logging.getLogger(__name__) router = APIRouter() def require_admin( x_admin_token: Annotated[str | None, Header()] = None, ) -> None: """Проверка admin-токена для /api/v1/admin/*. Если settings.admin_token не задан (dev) — пропускаем (открыто). Если задан — требуем заголовок X-Admin-Token с этим значением. """ if settings.admin_token is None: return # dev-режим: токен не настроен → эндпоинты открыты if x_admin_token != settings.admin_token: raise HTTPException(status_code=401, detail="Invalid or missing X-Admin-Token") class ScrapeRequest(BaseModel): lat: float lon: float radius_m: int = Field(default=1000, ge=100, le=20000) sources: list[Literal["avito", "cian", "yandex", "n1"]] = Field(default_factory=lambda: ["avito", "cian", "yandex", "n1"]) multi_room_yandex: bool = False # Если True — скрейп Yandex по 5 сегментам комнат, не одним общим запросом. deep_yandex: bool = False # Если True — Yandex полный обход 5 rooms × 3 sorts × 2 pages = 30 запросов / ~150s. multi_room_cian: bool = False # Если True — скрейп Cian по 4 сегментам комнат отдельно (~4× лотов). multi_room_n1: bool = False # Если True — скрейп N1 по 4 сегментам комнат отдельно (~4× лотов). 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, dependencies=[Depends(require_admin)]) 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] = [] for source in payload.sources: scraper_cls = { "avito": AvitoScraper, "cian": CianScraper, "yandex": YandexRealtyScraper, "n1": N1Scraper, }.get(source) if scraper_cls is None: continue async with scraper_cls() 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 ) elif source == "n1" and payload.multi_room_n1: 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 ) inserted, updated = save_listings(db, lots) 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. Берём часть до ' · '. N1 отдаёт «Репина, 75/2 стр.» — ок. """ main = addr.split(" · ")[0].strip() return main or addr @router.post("/geocode-missing", dependencies=[Depends(require_admin)]) async def geocode_missing( db: Annotated[Session, Depends(get_db)], limit: int = 100, ) -> dict: """Геокодинг listings у которых нет lat/lon (используя address). Чанк-обработка: limit=100 → ~150-200с (Nominatim 1 req/sec). cron вызывает в цикле пока `remaining` > 0. geocode_tried_at: после КАЖДОЙ попытки (успех/провал) ставим NOW(). Failed- адреса не выбираются повторно 7 дней → cron-loop завершается, не зацикливается. geom обновляется автоматически триггером listings_set_geom_trg. """ rows = db.execute( text( """ SELECT id, address FROM listings WHERE lat IS NULL AND COALESCE(address, '') != '' AND address NOT LIKE '%(Avito)%' AND address NOT LIKE '%(N1)%' AND (geocode_tried_at IS NULL OR geocode_tried_at < NOW() - interval '7 days') ORDER BY geocode_tried_at NULLS FIRST, scraped_at DESC LIMIT :limit """ ), {"limit": limit}, ).mappings().all() geocoded = 0 skipped = 0 for row in rows: clean = _clean_address_for_geocode(row["address"]) result = await geocode(clean, db) if result is None: # Помечаем что пробовали — иначе ретрай на каждом cron. db.execute( text("UPDATE listings SET geocode_tried_at = NOW() WHERE id = :id"), {"id": row["id"]}, ) db.commit() skipped += 1 continue # geom пересчитается триггером из lat/lon. db.execute( text( "UPDATE listings 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( """ SELECT count(*) FROM listings WHERE lat IS NULL AND COALESCE(address, '') != '' AND address NOT LIKE '%(Avito)%' AND address NOT LIKE '%(N1)%' 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), }