"""Загрузчик OSM POI из Overpass API для site-finder. Запускается раз в неделю через Celery beat. Поддерживает фильтр "не старше 2 лет" (требование Максима) — last_osm_edit_date. """ import asyncio import json import logging from datetime import date, datetime import httpx from sqlalchemy import text from app.core.db import SessionLocal logger = logging.getLogger(__name__) OVERPASS_URL = "https://overpass-api.de/api/interpreter" EKB_BBOX = (56.7, 60.5, 56.95, 60.75) # (south, west, north, east) # Маппинг OSM-тег → нормализованная category OSM_CATEGORIES: dict[tuple[str, str], str] = { # amenity tags ("amenity", "school"): "school", ("amenity", "kindergarten"): "kindergarten", ("amenity", "pharmacy"): "pharmacy", ("amenity", "hospital"): "hospital", ("amenity", "clinic"): "hospital", # shop tags ("shop", "mall"): "shop_mall", ("shop", "supermarket"): "shop_supermarket", ("shop", "convenience"): "shop_small", ("shop", "bakery"): "shop_small", # leisure ("leisure", "park"): "park", # transit (railway tram_stop приоритет над public_transport) ("railway", "tram_stop"): "tram_stop", ("highway", "bus_stop"): "bus_stop", # метро (одна линия в ЕКБ, но добавляем для полноты) ("station", "subway"): "metro_stop", } def _build_overpass_query() -> str: """Строит Overpass QL запрос для всех POI типов в bbox ЕКБ.""" south, west, north, east = EKB_BBOX bbox = f"({south},{west},{north},{east})" # out body даёт теги + meta для timestamp последней правки queries: list[str] = [] for (k, v), _cat in OSM_CATEGORIES.items(): queries.append(f'node["{k}"="{v}"]{bbox};') queries.append(f'way["{k}"="{v}"]{bbox};') inner = "\n".join(queries) return f"[out:json][timeout:60];(\n{inner}\n);out center meta;" def _classify(tags: dict[str, str]) -> str | None: """Определить category из OSM-тегов. None если не соответствует ни одной.""" for (k, v), cat in OSM_CATEGORIES.items(): if tags.get(k) == v: return cat return None async def fetch_overpass() -> list[dict]: """Запросить Overpass API, вернуть список raw OSM elements.""" query = _build_overpass_query() logger.info("Overpass: отправляем запрос (%d символов)", len(query)) async with httpx.AsyncClient(timeout=120) as client: r = await client.post(OVERPASS_URL, data={"data": query}) r.raise_for_status() elements: list[dict] = r.json().get("elements", []) logger.info("Overpass: получено %d элементов", len(elements)) return elements def sync_poi_to_db() -> dict[str, int]: """Синхронизирует POI из Overpass в osm_poi_ekb. UPSERT по UNIQUE(osm_type, osm_id, category). Returns: counters {fetched, inserted, updated, skipped_old}. """ elements = asyncio.run(fetch_overpass()) two_years_ago = date.today().replace(year=date.today().year - 2) inserted = 0 updated = 0 skipped_old = 0 fetched = len(elements) db = SessionLocal() try: for el in elements: tags: dict[str, str] = el.get("tags") or {} category = _classify(tags) if not category: continue osm_id: int = el["id"] osm_type: str = el["type"] # 'node' | 'way' | 'relation' # node имеет lat/lon напрямую, way — через center if osm_type == "node": lat = el.get("lat") lon = el.get("lon") else: center: dict = el.get("center") or {} lat = center.get("lat") lon = center.get("lon") if lat is None or lon is None: continue # timestamp последней правки элемента ts: str | None = el.get("timestamp") last_edit: date | None = None if ts: try: last_edit = datetime.fromisoformat(ts.replace("Z", "+00:00")).date() except Exception: last_edit = None # Мягкий фильтр "не старше 2 лет" — данные сохраняем, помечаем для API if last_edit and last_edit < two_years_ago: skipped_old += 1 result = db.execute( text(""" INSERT INTO osm_poi_ekb (osm_id, osm_type, category, name, lat, lon, geom, last_osm_edit_date, tags, fetched_at) VALUES (:osm_id, :osm_type, :category, :name, :lat, :lon, ST_SetSRID(ST_MakePoint(:lon, :lat), 4326), :last_edit, CAST(:tags AS jsonb), NOW()) ON CONFLICT (osm_type, osm_id, category) DO UPDATE SET name = EXCLUDED.name, lat = EXCLUDED.lat, lon = EXCLUDED.lon, geom = EXCLUDED.geom, last_osm_edit_date = EXCLUDED.last_osm_edit_date, tags = EXCLUDED.tags, fetched_at = NOW() RETURNING (xmax = 0) AS is_insert """), { "osm_id": osm_id, "osm_type": osm_type, "category": category, "name": tags.get("name"), "lat": lat, "lon": lon, "last_edit": last_edit, "tags": json.dumps(tags, ensure_ascii=False), }, ).scalar() if result: inserted += 1 else: updated += 1 db.commit() except Exception: db.rollback() raise finally: db.close() return { "fetched": fetched, "inserted": inserted, "updated": updated, "skipped_old": skipped_old, }