diff --git a/backend/app/services/site_finder/poi_loader.py b/backend/app/services/site_finder/poi_loader.py index d0473ad8..d4f7984e 100644 --- a/backend/app/services/site_finder/poi_loader.py +++ b/backend/app/services/site_finder/poi_loader.py @@ -1,12 +1,17 @@ """Загрузчик OSM POI из Overpass API для site-finder. -Запускается раз в неделю через Celery beat. Поддерживает фильтр -"не старше 2 лет" (требование Максима) — last_osm_edit_date. +Запускается раз в неделю через 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 @@ -19,6 +24,31 @@ 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 ниже) — это "ИЛИ" на @@ -54,17 +84,78 @@ OSM_CATEGORIES: dict[tuple[tuple[str, str], ...], str] = { } -def _build_overpass_query(tag_filters: tuple[tuple[str, str], ...]) -> str: - """Запрос для ОДНОЙ комбинации tag=value (обычно один тег, иногда несколько — все AND). +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 + + +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 = EKB_BBOX - bbox = f"({south},{west},{north},{east})" + 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};way{filt}{bbox};);out center meta;" + 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: @@ -80,25 +171,38 @@ def _tag_filters_desc(tag_filters: tuple[tuple[str, str], ...]) -> str: async def _fetch_category( - client: httpx.AsyncClient, tag_filters: tuple[tuple[str, str], ...], category: str + client: httpx.AsyncClient, + tag_filters: tuple[tuple[str, str], ...], + category: str, + bbox: tuple[float, float, float, float], + depth: int = 0, ) -> list[dict]: - """Один per-category Overpass-запрос с ОДНИМ повтором при транзиентной ошибке. + """Один per-category Overpass-запрос (для ОДНОГО тайла bbox) с ОДНИМ повтором при + транзиентной ошибке; если тайл падает оба раза — адаптивно дробится на 4 четверти + (см. RECURSIVE_SPLIT_MAX_DEPTH) и перезапрашивается рекурсивно, вместо того чтобы + тихо потерять весь тайл. Fix (location-index rework, "не потерялись крупные категории"): раньше единственная неудача (таймаут / 504) на всю неделю обнуляла категорию целиком (следующая попытка — только на следующем weekly run). Один retry с паузой снимает большую часть транзиентных сбоев без риска зациклиться (Overpass rate-limit — max 2 concurrent, поэтому не более - 2 попыток на категорию). + 2 попыток на тайл ДО дробления). """ tag_desc = _tag_filters_desc(tag_filters) - query = _build_overpass_query(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) → %d [attempt %d]", tag_desc, category, len(elements), attempt + "Overpass: %s (%s) bbox=%s depth=%d → %d [attempt %d]", + tag_desc, + category, + bbox, + depth, + len(elements), + attempt, ) # Привязываем category именно к тому per-category запросу, под который # элемент реально пришёл. Элемент с двумя целевыми тегами (например @@ -113,22 +217,49 @@ async def _fetch_category( logger.warning("Overpass failed for %s (attempt 1, retrying): %s", tag_desc, e) await asyncio.sleep(3.0) continue + if 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, depth + 1) + ) + await asyncio.sleep(1.0) + return combined logger.warning( - "Overpass failed for %s after retry — category skipped this run: %s", tag_desc, e + "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() -> list[dict]: - """Запросить Overpass API per category, вернуть combined список elements. +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). + 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) headers = { "User-Agent": "GenDesign-SiteFinder/1.0 (+https://gendsgn.ru)", "Accept": "application/json", @@ -136,24 +267,33 @@ async def fetch_overpass() -> list[dict]: 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) + for tile in tiles: + elements = await _fetch_category(client, tag_filters, category, tile) + all_elements.extend(elements) + await asyncio.sleep(1.0) logger.info( - "Overpass: total %d elements across %d category-queries", + "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() -> dict[str, int]: - """Синхронизирует POI из Overpass в osm_poi_ekb. +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), поэтому + НЕ переименована — см. ADR в vault decisions/. UPSERT по UNIQUE(osm_type, osm_id, category). Returns: counters {fetched, inserted, updated, skipped_old}. """ - elements = asyncio.run(fetch_overpass()) + 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) diff --git a/backend/tests/test_poi_loader.py b/backend/tests/test_poi_loader.py new file mode 100644 index 00000000..e1fb2515 --- /dev/null +++ b/backend/tests/test_poi_loader.py @@ -0,0 +1,178 @@ +"""Unit tests для poi_loader — региональный bbox + адаптивное дробление Overpass-тайлов. + +Mock-based / pure — БЕЗ живых походов в Overpass и без БД (правило: реальные запросы +к Overpass в тестах недопустимы). Покрывает: +- REGION_BBOX / DEFAULT_REGION — дефолт остаётся "ekb", не меняется молча. +- _bbox_tiles — для ЕКБ ровно один тайл, идентичный EKB_BBOX; для Москвы — несколько + тайлов, покрывающих исходный bbox без дыр/нахлёста (по площади). +- _build_overpass_query — bbox теперь параметр, а не глобальная константа. +- _fetch_category — retry (было и раньше) + НОВОЕ: адаптивное дробление тайла на 4 + четверти при устойчивом провале (вместо тихой потери тайла), с остановкой на + RECURSIVE_SPLIT_MAX_DEPTH (без бесконечной рекурсии). +""" + +from __future__ import annotations + +from types import SimpleNamespace + +import pytest + +from app.services.site_finder.poi_loader import ( + DEFAULT_REGION, + EKB_BBOX, + RECURSIVE_SPLIT_MAX_DEPTH, + REGION_BBOX, + _bbox_tiles, + _build_overpass_query, + _fetch_category, + _split_bbox_quadrants, +) + + +class _FakeResponse: + def __init__(self, ok: bool, elements: list[dict] | None = None) -> None: + self._ok = ok + self._elements = elements or [] + + def raise_for_status(self) -> None: + if not self._ok: + raise RuntimeError("simulated Overpass failure") + + def json(self) -> dict: + return {"elements": self._elements} + + +@pytest.fixture +def instant_sleep(monkeypatch: pytest.MonkeyPatch) -> None: + """Подменяет asyncio.sleep внутри poi_loader на no-op — тесты дробления тайлов иначе + реально спали бы минуты (retry-пауза 3с + 1с между каждой из 4 четвертей на каждом + уровне рекурсии).""" + + async def _instant(_seconds: float) -> None: + return None + + monkeypatch.setattr("app.services.site_finder.poi_loader.asyncio.sleep", _instant) + + +# ── REGION_BBOX / дефолт ────────────────────────────────────────────────────── + + +def test_default_region_is_ekb_unchanged() -> None: + assert DEFAULT_REGION == "ekb" + assert REGION_BBOX["ekb"] == EKB_BBOX + + +def test_region_bbox_has_msk_product_core() -> None: + assert "msk" in REGION_BBOX + south, west, north, east = REGION_BBOX["msk"] + assert south < north + assert west < east + + +# ── _bbox_tiles ──────────────────────────────────────────────────────────────── + + +def test_bbox_tiles_ekb_is_single_tile_identical_to_ekb_bbox() -> None: + """Дефолтный регион не должен молча поменять поведение — один тайл, байт-в-байт EKB_BBOX.""" + tiles = _bbox_tiles(EKB_BBOX) + assert tiles == [EKB_BBOX] + + +def test_bbox_tiles_msk_splits_into_multiple_tiles_without_gaps() -> None: + bbox = REGION_BBOX["msk"] + tiles = _bbox_tiles(bbox) + assert len(tiles) > 1 + south, west, north, east = bbox + total_area = (north - south) * (east - west) + tiles_area = sum((t[2] - t[0]) * (t[3] - t[1]) for t in tiles) + assert tiles_area == pytest.approx(total_area, rel=1e-9) + + +# ── _split_bbox_quadrants ──────────────────────────────────────────────────────── + + +def test_split_bbox_quadrants_covers_original_area() -> None: + bbox = (55.55, 37.30, 55.95, 37.90) + quads = _split_bbox_quadrants(bbox) + assert len(quads) == 4 + south, west, north, east = bbox + total_area = (north - south) * (east - west) + quads_area = sum((q[2] - q[0]) * (q[3] - q[1]) for q in quads) + assert quads_area == pytest.approx(total_area, rel=1e-9) + + +# ── _build_overpass_query ───────────────────────────────────────────────────── + + +def test_build_overpass_query_uses_given_bbox_not_global_constant() -> None: + q = _build_overpass_query((("amenity", "pharmacy"),), (1.0, 2.0, 3.0, 4.0)) + assert "(1.0,2.0,3.0,4.0)" in q + assert '["amenity"="pharmacy"]' in q + + +# ── _fetch_category: retry (существующее поведение) ────────────────────────────── + + +async def test_fetch_category_retries_then_succeeds(instant_sleep: None) -> None: + calls = {"n": 0} + + async def fake_post(_url: str, data: dict) -> _FakeResponse: + calls["n"] += 1 + if calls["n"] == 1: + return _FakeResponse(ok=False) + return _FakeResponse(ok=True, elements=[{"type": "node", "id": 1, "lat": 1, "lon": 2}]) + + client = SimpleNamespace(post=fake_post) + result = await _fetch_category(client, (("amenity", "pharmacy"),), "pharmacy", (0, 0, 1, 1)) + assert calls["n"] == 2 + assert len(result) == 1 + assert result[0]["_gd_category"] == "pharmacy" + + +# ── _fetch_category: адаптивное дробление (НОВОЕ) ───────────────────────────────── + + +async def test_fetch_category_splits_into_quadrants_on_persistent_failure( + instant_sleep: None, +) -> None: + """Тайл, где оба attempt проваливаются, дробится на 4 четверти вместо потери данных.""" + calls = {"n": 0} + + async def fake_post(_url: str, data: dict) -> _FakeResponse: + calls["n"] += 1 + query = data["data"] + if "(0.0,0.0,1.0,1.0)" in query: # верхнеуровневый тайл всегда 504 + return _FakeResponse(ok=False) + return _FakeResponse( + ok=True, elements=[{"type": "node", "id": calls["n"], "lat": 0.1, "lon": 0.1}] + ) + + client = SimpleNamespace(post=fake_post) + result = await _fetch_category( + client, (("amenity", "pharmacy"),), "pharmacy", (0.0, 0.0, 1.0, 1.0) + ) + # верхний тайл: 2 неудачных attempt, затем 4 успешных запроса по четвертям + assert calls["n"] == 2 + 4 + assert len(result) == 4 + + +async def test_fetch_category_gives_up_at_max_depth_without_infinite_recursion( + instant_sleep: None, +) -> None: + """Тайл, падающий на ЛЮБОМ размере, останавливает дробление на RECURSIVE_SPLIT_MAX_DEPTH + и возвращает пустой список — не зацикливается и не падает.""" + calls = {"n": 0} + + async def fake_post(_url: str, data: dict) -> _FakeResponse: + calls["n"] += 1 + assert data # параметр используется — сигнатура должна совпадать с client.post + return _FakeResponse(ok=False) + + client = SimpleNamespace(post=fake_post) + result = await _fetch_category( + client, (("amenity", "pharmacy"),), "pharmacy", (0.0, 0.0, 1.0, 1.0) + ) + assert result == [] + # sum_{d=0}^{max_depth} 4^d узлов, каждый по 2 attempt — рекурсия конечна + expected_nodes = sum(4**d for d in range(RECURSIVE_SPLIT_MAX_DEPTH + 1)) + assert calls["n"] == expected_nodes * 2