"""Загрузчик OSM POI из Overpass API для site-finder. Запускается раз в неделю через Celery beat. Поддерживает фильтр "не старше 2 лет" (требование Максима) — last_osm_edit_date. """ import asyncio import json import logging from datetime import date, datetime, timedelta 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-тегов (все теги в кортеже должны совпасть — AND) → нормализованная # category. Каждая запись — один per-category Overpass-запрос (см. _build_overpass_query); # несколько записей с ОДИНАКОВЫМ значением category (как у metro_stop ниже) — это "ИЛИ" на # уровне отдельных HTTP-запросов: элемент, подходящий под любую из альтернативных схем # разметки, попадёт в категорию. OSM_CATEGORIES: dict[tuple[tuple[str, str], ...], str] = { # amenity tags — школы расширены (school/college/university) (("amenity", "school"),): "school", (("amenity", "college"),): "school", (("amenity", "university"),): "school", (("amenity", "kindergarten"),): "kindergarten", (("amenity", "pharmacy"),): "pharmacy", (("amenity", "hospital"),): "hospital", (("amenity", "clinic"),): "hospital", # shop tags — supermarket расширен (("shop", "mall"),): "shop_mall", (("shop", "supermarket"),): "shop_supermarket", (("shop", "hypermarket"),): "shop_supermarket", (("shop", "convenience"),): "shop_small", (("shop", "bakery"),): "shop_small", # leisure (("leisure", "park"),): "park", # transit (("railway", "tram_stop"),): "tram_stop", (("highway", "bus_stop"),): "bus_stop", # Метро ЕКБ (9 станций, одна линия). Fix (location-index rework): фильтр раньше ловил # ТОЛЬКО station=subway и подтягивал лишь 5/9 станций — часть станций в OSM размечена # без ключа "station" вовсе, комбинацией railway=station + subway=yes (альтернативная, # но распространённая схема разметки метро). Обе схемы — отдельными записями ниже, чтобы # не терять станции, размеченные любой из них. (("station", "subway"),): "metro_stop", (("railway", "station"), ("subway", "yes")): "metro_stop", } def _build_overpass_query(tag_filters: tuple[tuple[str, str], ...]) -> str: """Запрос для ОДНОЙ комбинации tag=value (обычно один тег, иногда несколько — все AND). Раньше делали один большой запрос на все 14 категорий — Overpass возвращал 504 Gateway Timeout (запрос слишком тяжёлый). Сплит на per-category даёт быстрые запросы вместо одного 60+ секундного. """ south, west, north, east = EKB_BBOX bbox = f"({south},{west},{north},{east})" filt = "".join(f'["{k}"="{v}"]' for k, v in tag_filters) return f"[out:json][timeout:30];(node{filt}{bbox};way{filt}{bbox};);out center meta;" def _classify(tags: dict[str, str]) -> str | None: """Определить category из OSM-тегов. None если не соответствует ни одной.""" for tag_filters, cat in OSM_CATEGORIES.items(): if all(tags.get(k) == v for k, v in tag_filters): return cat return None def _tag_filters_desc(tag_filters: tuple[tuple[str, str], ...]) -> str: return ",".join(f"{k}={v}" for k, v in tag_filters) async def _fetch_category( client: httpx.AsyncClient, tag_filters: tuple[tuple[str, str], ...], category: str ) -> list[dict]: """Один per-category Overpass-запрос с ОДНИМ повтором при транзиентной ошибке. Fix (location-index rework, "не потерялись крупные категории"): раньше единственная неудача (таймаут / 504) на всю неделю обнуляла категорию целиком (следующая попытка — только на следующем weekly run). Один retry с паузой снимает большую часть транзиентных сбоев без риска зациклиться (Overpass rate-limit — max 2 concurrent, поэтому не более 2 попыток на категорию). """ tag_desc = _tag_filters_desc(tag_filters) query = _build_overpass_query(tag_filters) for attempt in (1, 2): try: r = await client.post(OVERPASS_URL, data={"data": query}) r.raise_for_status() elements: list[dict] = r.json().get("elements", []) logger.info( "Overpass: %s (%s) → %d [attempt %d]", tag_desc, category, len(elements), attempt ) # Привязываем category именно к тому per-category запросу, под который # элемент реально пришёл. Элемент с двумя целевыми тегами (например # amenity=pharmacy + shop=supermarket) приходит дважды — каждая копия # несёт свою category. Иначе _classify по dict-порядку молча терял бы # вторую категорию при UPSERT по UNIQUE(osm_type, osm_id, category). См. #1372. for el in elements: el["_gd_category"] = category return elements except Exception as e: if attempt == 1: logger.warning("Overpass failed for %s (attempt 1, retrying): %s", tag_desc, e) await asyncio.sleep(3.0) continue logger.warning( "Overpass failed for %s after retry — category skipped this run: %s", tag_desc, e ) return [] async def fetch_overpass() -> list[dict]: """Запросить Overpass API per category, вернуть combined список elements. Делаем отдельные запросы вместо одного гигантского — большой запрос отдаёт 504 Gateway Timeout. Между запросами sleep 1с (Overpass usage policy: max 2 concurrent, лучше 1 req/s). Overpass блокирует default `python-httpx/*` User-Agent (406) — поэтому явный UA с контактом проекта. """ headers = { "User-Agent": "GenDesign-SiteFinder/1.0 (+https://gendsgn.ru)", "Accept": "application/json", } all_elements: list[dict] = [] async with httpx.AsyncClient(timeout=60, headers=headers) as client: for tag_filters, category in OSM_CATEGORIES.items(): elements = await _fetch_category(client, tag_filters, category) all_elements.extend(elements) await asyncio.sleep(1.0) logger.info( "Overpass: total %d elements across %d category-queries", len(all_elements), len(OSM_CATEGORIES), ) return all_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()) # 730 дней ≈ 2 года: избегаем ValueError 29 февраля (year-2 не високосный → нет 29.02). # Точность ±1 день несущественна для фильтра "не старше 2 лет" (требование Максима). См. #1232. two_years_ago = date.today() - timedelta(days=730) inserted = 0 updated = 0 skipped_old = 0 skipped = 0 fetched = len(elements) db = SessionLocal() try: for el in elements: tags: dict[str, str] = el.get("tags") or {} # category проставлена в fetch_overpass под тот per-category запрос, под # который элемент пришёл (multi-tag элемент приходит несколько раз, каждая # копия со своей category). Fallback на _classify для прямых вызовов. См. #1372. category = el.get("_gd_category") or _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 try: with db.begin_nested(): # SAVEPOINT — откат только этой записи 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 except Exception as e: # Дефектный OSM-элемент (битый tags-jsonb, нарушение constraint, PostGIS- # ошибка у way/relation без корректного center) не должен валить весь # weekly-sync. SAVEPOINT откатывает только эту строку, продолжаем # с остальными. См. #1343 и .claude/rules/backend.md (SAVEPOINT pattern). logger.warning( "poi_sync insert failed for %s/%s (category=%s): %s", osm_type, osm_id, category, e, ) skipped += 1 db.commit() except Exception as e: db.rollback() logger.exception("poi_sync: unexpected error, outer tx rolled back: %s", e) raise finally: db.close() return { "fetched": fetched, "inserted": inserted, "updated": updated, "skipped_old": skipped_old, "skipped": skipped, }