"""Загрузчик OSM POI из Overpass API для site-finder. Запускается раз в неделю через Celery beat (регион по умолчанию — ЕКБ, см. DEFAULT_REGION). Поддерживает фильтр "не старше 2 лет" (требование Максима) — last_osm_edit_date. Параметризован регионом (REGION_BBOX) — sync_poi_to_db(region=...) может грузить любой зарегистрированный bbox, не только ЕКБ; большие bbox автоматически режутся на тайлы (_bbox_tiles), чтобы не упереться в лимиты одного Overpass-запроса. """ import asyncio import json import logging import math 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) # Прямоугольники продуктовых ядер по региону — источник загрузки POI больше не зашит # в одну константу (было: только ЕКБ, блок «что рядом» молчал для остальных регионов). # Site Finder — независимая половина монорепо со своим окружением/БД и НЕ импортирует # tradein-mvp (у того свой реестр `app/services/regions.py`), поэтому bbox для Москвы # продублирован явно, а не через кросс-импорт. Значение — bbox_product_core региона 77 # (tradein-mvp/backend/app/services/regions.py REGIONS[77]), пересчитанное в тот же # (south, west, north, east) порядок, что EKB_BBOX выше. "ekb" остаётся значением по # умолчанию ВЕЗДЕ (sync_poi_to_db / fetch_overpass) — существующее weekly-расписание # (tasks.poi_sync.sync_osm_poi_ekb) не передаёт region и не должно молча сменить город. REGION_BBOX: dict[str, tuple[float, float, float, float]] = { "ekb": EKB_BBOX, "msk": (55.55, 37.30, 55.95, 37.90), # (south, west, north, east) } DEFAULT_REGION = "ekb" # Максимальный размер стороны ОДНОГО Overpass-запроса в градусах. У ЕКБ обе стороны # bbox — ровно 0.25° (проверенный на практике размер: per-category запрос укладывается # в timeout:30 без 504). Для региона с большей стороной bbox запрос режется на грид # тайлов такого же порядка вместо одного большого — иначе на плотном городе (Москва на # порядок плотнее ЕКБ по числу POI) Overpass либо отдаёт 504, либо (хуже) частично # посчитанный ответ без явной ошибки, и загрузка молча обрежется. Для ЕКБ (0.25×0.25) # тайлинг даёт РОВНО один тайл, совпадающий с EKB_BBOX бит-в-бит — поведение дефолтного # региона не меняется. MAX_TILE_SIDE_DEG = 0.25 # Маппинг набора 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 _bbox_tiles( bbox: tuple[float, float, float, float], max_side_deg: float = MAX_TILE_SIDE_DEG ) -> list[tuple[float, float, float, float]]: """Разбить bbox на равномерный грид тайлов со стороной ≤ max_side_deg. (south, west, north, east) → список тайлов того же формата. Для bbox, у которого обе стороны уже ≤ max_side_deg (текущий EKB_BBOX: 0.25×0.25), возвращает список ровно из ОДНОГО тайла, идентичного входному bbox — тайлинг не меняет поведение для ЕКБ. """ south, west, north, east = bbox rows = max(1, math.ceil(round((north - south) / max_side_deg, 6))) cols = max(1, math.ceil(round((east - west) / max_side_deg, 6))) lat_step = (north - south) / rows lon_step = (east - west) / cols tiles = [] for r in range(rows): for c in range(cols): tiles.append( ( south + r * lat_step, west + c * lon_step, south + (r + 1) * lat_step, west + (c + 1) * lon_step, ) ) return tiles def _split_bbox_quadrants( bbox: tuple[float, float, float, float], ) -> list[tuple[float, float, float, float]]: """Разбить bbox на 4 равные четверти (2×2) — используется адаптивным ретраем _fetch_category, когда сам тайл всё равно оказался слишком тяжёлым для Overpass.""" south, west, north, east = bbox mid_lat = (south + north) / 2 mid_lon = (west + east) / 2 return [ (south, west, mid_lat, mid_lon), (south, mid_lon, mid_lat, east), (mid_lat, west, north, mid_lon), (mid_lat, mid_lon, north, east), ] # Живой замер 2026-09-13: uniform-тайл 0.25×0.25 (размер ЕКБ) для category=bus_stop в # Москве отдал 504 Gateway Timeout ОДНИМ тайлом (12 181 bus_stop во всём продуктовом # ядре — уже больше, чем ВСЕ 14 категорий ЕКБ вместе, 4 850). Единый "правильный" размер # тайла под все 14 категорий Москвы заранее не подобрать — плотность по городу сильно # неравномерна (плотный центр / разреженная периферия), а у Overpass нет заголовка с # "это частичный ответ" — единственный надёжный сигнал перегруза — HTTP-ошибка/таймаут. # Поэтому вместо фиксированного маленького тайла — АДАПТИВНОЕ дробление: тайл, на # котором per-category запрос дважды падает, дробится на 4 четверти и каждая # перезапрашивается рекурсивно (до RECURSIVE_SPLIT_MAX_DEPTH). Для ЕКБ recursion # НИКОГДА не срабатывает (единственный тайл исторически всегда отвечал 200) — поведение # дефолтного региона не меняется. RECURSIVE_SPLIT_MAX_DEPTH = 3 # Дробить тайл имеет смысл ТОЛЬКО когда сервер отказал из-за тяжести запроса: # 504/429/503 и таймаут чтения — «я не успел посчитать», четверть посчитается. # Отказ на уровне транспорта (connection refused / network unreachable) про размер # запроса не говорит ВООБЩЕ: хост нас не принимает, и дробление превращает один # отказ в 4, 16, 64 повторных стука. Живой случай 15.09.2026: загрузка Москвы # поймала блокировку overpass-api.de по IP и за три минуты выдала 58 отказов на # 4 успеха — ровно этот механизм. _SPLIT_WORTHY_STATUS = frozenset({429, 503, 504}) # Подряд идущие транспортные отказы = хост нас не принимает. Продолжать прогон # бессмысленно и вредно (углубляем блокировку), поэтому после порога — стоп всего # прогона с явной ошибкой, а не тихий пропуск категорий. _MAX_CONSECUTIVE_TRANSPORT_ERRORS = 5 class OverpassUnreachableError(RuntimeError): """Overpass отказывает на уровне соединения подряд — прогон остановлен.""" class _RunState: """Счётчик подряд идущих транспортных отказов в рамках одного fetch_overpass.""" __slots__ = ("consecutive_transport_errors",) def __init__(self) -> None: self.consecutive_transport_errors = 0 def _is_overload(exc: Exception) -> bool: """True, если отказ говорит «запрос слишком тяжёлый» (есть смысл дробить).""" if isinstance(exc, httpx.TimeoutException): return True if isinstance(exc, httpx.HTTPStatusError): return exc.response.status_code in _SPLIT_WORTHY_STATUS return False def _is_transport_error(exc: Exception) -> bool: """True для отказа на уровне соединения (хост не принимает), не про размер запроса.""" return isinstance(exc, httpx.TransportError) and not isinstance(exc, httpx.TimeoutException) def _build_overpass_query( tag_filters: tuple[tuple[str, str], ...], bbox: tuple[float, float, float, float] ) -> str: """Запрос для ОДНОЙ комбинации tag=value (обычно один тег, иногда несколько — все AND) в ОДНОМ тайле bbox (south, west, north, east) — см. _bbox_tiles. Раньше делали один большой запрос на все 14 категорий — Overpass возвращал 504 Gateway Timeout (запрос слишком тяжёлый). Сплит на per-category даёт быстрые запросы вместо одного 60+ секундного. """ south, west, north, east = bbox bbox_str = 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_str};way{filt}{bbox_str};);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, bbox: tuple[float, float, float, float], state: _RunState, depth: int = 0, ) -> list[dict]: """Один per-category Overpass-запрос (для ОДНОГО тайла bbox) с ОДНИМ повтором при транзиентной ошибке; если тайл падает оба раза — адаптивно дробится на 4 четверти (см. RECURSIVE_SPLIT_MAX_DEPTH) и перезапрашивается рекурсивно, вместо того чтобы тихо потерять весь тайл. 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, bbox) 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) bbox=%s depth=%d → %d [attempt %d]", tag_desc, category, bbox, depth, len(elements), attempt, ) # Привязываем category именно к тому per-category запросу, под который # элемент реально пришёл. Элемент с двумя целевыми тегами (например # amenity=pharmacy + shop=supermarket) приходит дважды — каждая копия # несёт свою category. Иначе _classify по dict-порядку молча терял бы # вторую категорию при UPSERT по UNIQUE(osm_type, osm_id, category). См. #1372. state.consecutive_transport_errors = 0 for el in elements: el["_gd_category"] = category return elements except Exception as e: if _is_transport_error(e): state.consecutive_transport_errors += 1 if state.consecutive_transport_errors >= _MAX_CONSECUTIVE_TRANSPORT_ERRORS: raise OverpassUnreachableError( f"Overpass отказывает на уровне соединения " f"{state.consecutive_transport_errors} раз подряд ({e}) — прогон " f"остановлен, чтобы не стучаться в блокирующий хост" ) from e logger.warning( "Overpass transport error for %s bbox=%s (подряд %d) — тайл пропущен " "без дробления: %s", tag_desc, bbox, state.consecutive_transport_errors, e, ) return [] state.consecutive_transport_errors = 0 if attempt == 1: logger.warning("Overpass failed for %s (attempt 1, retrying): %s", tag_desc, e) await asyncio.sleep(3.0) continue if _is_overload(e) and depth < RECURSIVE_SPLIT_MAX_DEPTH: logger.warning( "Overpass failed for %s bbox=%s twice — splitting into 4 quadrants " "(depth %d→%d) instead of dropping the tile: %s", tag_desc, bbox, depth, depth + 1, e, ) combined: list[dict] = [] for quadrant in _split_bbox_quadrants(bbox): combined.extend( await _fetch_category( client, tag_filters, category, quadrant, state, depth + 1 ) ) await asyncio.sleep(1.0) return combined logger.warning( "Overpass failed for %s bbox=%s at max split depth %d — tile skipped this run: %s", tag_desc, bbox, depth, e, ) return [] async def fetch_overpass(region: str = DEFAULT_REGION) -> list[dict]: """Запросить Overpass API per category × per tile, вернуть combined список elements. Делаем отдельные запросы вместо одного гигантского — большой запрос отдаёт 504 Gateway Timeout. Между запросами sleep 1с (Overpass usage policy: max 2 concurrent, лучше 1 req/s). bbox региона режется на тайлы ≤ MAX_TILE_SIDE_DEG (_bbox_tiles) — для "ekb" это ровно один тайл (без изменения поведения), для регионов с большим bbox (напр. "msk") — несколько, чтобы не поймать 504 или тихо обрезанный ответ на плотном городе. Overpass блокирует default `python-httpx/*` User-Agent (406) — поэтому явный UA с контактом проекта. """ bbox = REGION_BBOX[region] tiles = _bbox_tiles(bbox) state = _RunState() 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(): for tile in tiles: elements = await _fetch_category(client, tag_filters, category, tile, state) all_elements.extend(elements) await asyncio.sleep(1.0) logger.info( "Overpass region=%s: total %d elements across %d category-queries × %d tiles", region, len(all_elements), len(OSM_CATEGORIES), len(tiles), ) return all_elements def sync_poi_to_db(region: str = DEFAULT_REGION) -> dict[str, int]: """Синхронизирует POI из Overpass в osm_poi_ekb для одного региона. region — ключ REGION_BBOX ("ekb" по умолчанию, сохраняет старое поведение weekly-расписания). Имя таблицы osm_poi_ekb — историческое (изначально ЕКБ-only); таблица читается ещё в двух местах вне Site Finder (FDW-таблица gendesign_osm_poi_ekb + локальное зеркало osm_poi_ekb_local в tradein), поэтому НЕ переименована: переименование потянуло бы миграции в обеих половинах монорепо (FDW-объект + зеркало + их индексы) ради косметики. UPSERT по UNIQUE(osm_type, osm_id, category). Returns: counters {fetched, inserted, updated, skipped_old}. """ elements = asyncio.run(fetch_overpass(region)) # 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, }