feat(site-finder): poi_loader грузит POI не только по Екатеринбургу

Блок «что рядом» в оценке молчал для Москвы — poi_loader.py тянул Overpass
строго по EKB_BBOX константе. Параметризовал регионом: REGION_BBOX содержит
"ekb" (дефолт, поведение weekly-расписания tasks.poi_sync.sync_osm_poi_ekb не
меняется) и "msk" (bbox_product_core региона 77 из tradein-mvp regions.py,
продублирован явно — Site Finder не импортирует tradein-mvp).

Живой замер Overpass (13.09.2026, count-only + один тайл) показал, что просто
разрешить region было бы недостаточно: в Москве 12 181 bus_stop — больше, чем
ВСЕ 14 категорий ЕКБ вместе (4 850) — и уже один uniform-тайл 0.25×0.25° (размер
ЕКБ) вернул 504 на этой категории. Единый безопасный размер тайла заранее не
подобрать — плотность Москвы неравномерна (плотный центр / разреженная
периферия), а у Overpass нет сигнала "это частичный ответ", кроме HTTP-ошибки.
Поэтому вместо фиксированного тайла — адаптивное дробление: тайл, дважды
упавший в Overpass, дробится на 4 четверти и перезапрашивается рекурсивно (до
RECURSIVE_SPLIT_MAX_DEPTH=3), вместо того чтобы тихо потерять данные по тайлу.
Для ЕКБ (единственный тайл исторически всегда отвечал 200) рекурсия не
срабатывает — поведение дефолта не меняется.

Таблицу osm_poi_ekb НЕ переименовывал: она читается ещё в двух местах вне
Site Finder (FDW gendesign_osm_poi_ekb + локальное зеркало osm_poi_ekb_local
в tradein), переименование задело бы миграции обеих половин ради красоты
имени. tradein-mvp/backend/app/tasks/osm_poi_ekb_refresh.py менять не
потребовалось — зеркало уже делает безусловный full-scan FDW-таблицы без
фильтра по bbox, так что московские точки подхватятся тем же кодом.

Загрузку Москвы в прод не запускал — по заданию это делает главная сессия:
sync_poi_to_db(region="msk").
This commit is contained in:
bot-backend 2026-09-13 15:44:35 +03:00
parent 811b9ecc8f
commit 33b72bac2e
2 changed files with 341 additions and 23 deletions

View file

@ -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)

View file

@ -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