feat(site-finder): backfill ЕКБ зон-регламента — разблокировать видимость финмодели (#1881)
financial_estimate (эпик #1881) фаерит только при разрезолвленном НСПД-регламенте (max_far), а он резолвится ЛЕНИВО при анализе участка → zone_regulation_cache всего 33 зоны → у большинства участков financial_estimate=None («—» в кокпите). Backfill проактивно кэширует ВСЕ террзоны ЕКБ (~100-120) → ЛЮБОй участок в известной зоне сразу получает регламент и финмодель. - backfill_ekb_zone_regulations (zone_regulation.py): WFS-перечисление террзон ЕКБ (features_in_bbox 'territorial_zone') → dedup по urban_index (предпочёт фичу с геометрией) → representative_point (внутри полигона, не centroid) → zone_regulation_at → upsert. ИДЕМПОТЕНТНО (cached skip перед fetch), per-zone try/except (один сбой не рушит batch), rate_delay вежливость к геопорталу, limit для частичного прогона. - Celery task backfill_zone_regulations + регистрация в celery_app include + admin POST /api/v1/admin/scrape/zone-regulations/backfill (паттерн objective sync-our). - Ленивый refresh_zone_regulations / get_or_fetch / upsert НЕ тронуты (additive путь). - +9 тестов (dedup, idempotency cached-skip, per-zone error isolation, centroid-inside, limit). mypy/ruff clean. Прод-прогон backfill — отдельный шаг после deploy (как Objective sync), не на deploy. Code-review поймал незарегистрированный task (был бы NotRegistered) — исправлено. Refs #1881
This commit is contained in:
parent
b7629067cf
commit
2ff34b1500
6 changed files with 495 additions and 0 deletions
|
|
@ -435,6 +435,54 @@ def trigger_pzz_sync() -> dict[str, Any]:
|
||||||
return {"task_id": result.id, "queued_at": "now"}
|
return {"task_id": result.id, "queued_at": "now"}
|
||||||
|
|
||||||
|
|
||||||
|
class TriggerZoneRegulationsBackfillRequest(BaseModel):
|
||||||
|
"""Параметры ручного backfill ПЗЗ-регламента терзон ЕКБ (эпик #1881)."""
|
||||||
|
|
||||||
|
limit: int | None = Field(
|
||||||
|
default=None,
|
||||||
|
ge=1,
|
||||||
|
le=1000,
|
||||||
|
description="Максимум уникальных незакэшированных зон к фетчу (None = все). Smoke-тест.",
|
||||||
|
)
|
||||||
|
rate_delay_s: float = Field(
|
||||||
|
default=1.0,
|
||||||
|
ge=0.0,
|
||||||
|
le=10.0,
|
||||||
|
description="Пауза (сек) между live-fetch'ами зон — вежливость к геопорталу.",
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@router.post("/zone-regulations/backfill")
|
||||||
|
def trigger_zone_regulations_backfill(
|
||||||
|
payload: TriggerZoneRegulationsBackfillRequest,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
"""Backfill кэша ПЗЗ-регламента ВСЕХ терзон ЕКБ → zone_regulation_cache (эпик #1881).
|
||||||
|
|
||||||
|
Чтобы financial_estimate фаерил для большинства участков (max_far из кэша, без live-urbanCard
|
||||||
|
в hot-пути). Идемпотентно: уже закэшированные зоны пропускаются, повторный запуск дёшев.
|
||||||
|
|
||||||
|
estimate_minutes — груба: ~120 уникальных терзон ЕКБ × (rate_delay + ~2с live-urbanCard).
|
||||||
|
"""
|
||||||
|
from app.workers.tasks.backfill_zone_regulations import backfill_zone_regulations
|
||||||
|
|
||||||
|
result = backfill_zone_regulations.apply_async(
|
||||||
|
kwargs={
|
||||||
|
"limit": payload.limit,
|
||||||
|
"rate_delay_s": payload.rate_delay_s,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
# ~120 уникальных терзон ЕКБ; на зону ≈ rate_delay + ~2с urbanCard round-trip.
|
||||||
|
n_zones = payload.limit if payload.limit is not None else 120
|
||||||
|
estimate_minutes = round(n_zones * (payload.rate_delay_s + 2.0) / 60.0, 1)
|
||||||
|
return {
|
||||||
|
"task_id": result.id,
|
||||||
|
"limit": payload.limit,
|
||||||
|
"rate_delay_s": payload.rate_delay_s,
|
||||||
|
"estimate_minutes": estimate_minutes,
|
||||||
|
"queued_at": "now",
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
@router.post("/poi-sync")
|
@router.post("/poi-sync")
|
||||||
def trigger_poi_sync() -> dict[str, Any]:
|
def trigger_poi_sync() -> dict[str, Any]:
|
||||||
"""Manual trigger для синхронизации OSM POI ЕКБ из Overpass API.
|
"""Manual trigger для синхронизации OSM POI ЕКБ из Overpass API.
|
||||||
|
|
|
||||||
|
|
@ -28,19 +28,27 @@ from __future__ import annotations
|
||||||
import json
|
import json
|
||||||
import logging
|
import logging
|
||||||
import re
|
import re
|
||||||
|
import time
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
|
from shapely.geometry import shape
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
from sqlalchemy.exc import OperationalError, ProgrammingError
|
from sqlalchemy.exc import OperationalError, ProgrammingError
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
from app.services.scrapers.ekb_geoportal_client import (
|
from app.services.scrapers.ekb_geoportal_client import (
|
||||||
|
EKBFeature,
|
||||||
EKBGeoportalClient,
|
EKBGeoportalClient,
|
||||||
ZoneRegulation,
|
ZoneRegulation,
|
||||||
)
|
)
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
# bbox агломерации ЕКБ (minlon, minlat, maxlon, maxlat) EPSG:4326 — границы для WFS-перечисления
|
||||||
|
# терзон. Слой territorial_zone геопортала — ЕКБ-only, поэтому bbox это просто bounds-обёртка
|
||||||
|
# (а не фильтр города). Шире чем нужно — лучше захватить лишнее, дедуп по zone_index снимет повтор.
|
||||||
|
_EKB_CITY_BBOX: tuple[float, float, float, float] = (60.0, 56.6, 61.1, 57.1)
|
||||||
|
|
||||||
# ── Числовой экстракт предельных параметров (regex) ──────────────────────────
|
# ── Числовой экстракт предельных параметров (regex) ──────────────────────────
|
||||||
# Числовой токен: тысячные разделены пробелом/nbsp ('50 000'), десятичная — запятая ('2,4').
|
# Числовой токен: тысячные разделены пробелом/nbsp ('50 000'), десятичная — запятая ('2,4').
|
||||||
# Сначала grouped-тысячи, иначе простое число. `.` по умолчанию НЕ переходит \n — параметры
|
# Сначала grouped-тысячи, иначе простое число. `.` по умолчанию НЕ переходит \n — параметры
|
||||||
|
|
@ -275,7 +283,136 @@ def get_or_fetch_zone_regulation(
|
||||||
return get_cached_zone_regulation(db, reg.zone_index, city=city)
|
return get_cached_zone_regulation(db, reg.zone_index, city=city)
|
||||||
|
|
||||||
|
|
||||||
|
# ── Backfill: проактивный прогрев ВСЕХ терзон ЕКБ (эпик #1881) ────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def _representative_point(geometry: dict[str, Any] | None) -> tuple[float, float] | None:
|
||||||
|
"""Точка (lon, lat) гарантированно ВНУТРИ полигона зоны.
|
||||||
|
|
||||||
|
``representative_point`` (не ``centroid``) гарантирует попадание внутрь даже у вогнутых
|
||||||
|
полигонов → searchByGeom/zone_regulation_at стабильно находит терзону. None — битая
|
||||||
|
геометрия (фичу пропускаем).
|
||||||
|
"""
|
||||||
|
if not geometry:
|
||||||
|
return None
|
||||||
|
try:
|
||||||
|
pt = shape(geometry).representative_point()
|
||||||
|
return (float(pt.x), float(pt.y))
|
||||||
|
except Exception as exc: # битая геометрия фичи — пропускаем зону, не валим batch
|
||||||
|
logger.warning("backfill_ekb_zone_regulations: bad geometry, skip: %s", exc)
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _dedup_zones_by_index(features: list[EKBFeature]) -> dict[str, EKBFeature]:
|
||||||
|
"""Сгруппировать фичи по ``urban_index``, выбрав одну фичу С ВАЛИДНОЙ геометрией на индекс.
|
||||||
|
|
||||||
|
Регламент идентичен для всех полигонов одной терзоны → резолвим индекс один раз.
|
||||||
|
Предпочитаем фичу с непустой геометрией (первая встретившаяся побеждает); фича без
|
||||||
|
геометрии берётся только если иной для индекса нет (всё равно отсеется в backfill).
|
||||||
|
"""
|
||||||
|
by_index: dict[str, EKBFeature] = {}
|
||||||
|
for feat in features:
|
||||||
|
idx = feat.properties.get("urban_index")
|
||||||
|
if not idx:
|
||||||
|
continue
|
||||||
|
key = str(idx)
|
||||||
|
existing = by_index.get(key)
|
||||||
|
if existing is None:
|
||||||
|
by_index[key] = feat
|
||||||
|
elif existing.geometry is None and feat.geometry is not None:
|
||||||
|
by_index[key] = feat # апгрейд до фичи с валидной геометрией
|
||||||
|
return by_index
|
||||||
|
|
||||||
|
|
||||||
|
def backfill_ekb_zone_regulations(
|
||||||
|
db: Session,
|
||||||
|
*,
|
||||||
|
client: EKBGeoportalClient | None = None,
|
||||||
|
bbox: tuple[float, float, float, float] = _EKB_CITY_BBOX,
|
||||||
|
limit: int | None = None,
|
||||||
|
rate_delay_s: float = 1.0,
|
||||||
|
city: str = "ekb",
|
||||||
|
) -> dict[str, int]:
|
||||||
|
"""Закэшировать регламент ВСЕХ терзон ЕКБ в ``zone_regulation_cache`` (эпик #1881).
|
||||||
|
|
||||||
|
Перечисляет терзоны через WFS, дедуплицирует по ``urban_index`` и для каждого УНИКАЛЬНОГО
|
||||||
|
индекса, которого ещё нет в кэше, резолвит регламент (центроид полигона → urbanCard) и
|
||||||
|
апсёртит. Идемпотентно: уже закэшированные зоны пропускаются (cached skip), повторный
|
||||||
|
прогон ничего не дофетчит. Per-zone try/except — сбой одной зоны не валит batch.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
db: SQLAlchemy-сессия (sync).
|
||||||
|
client: WFS-клиент геопортала; None → дефолтный ``EKBGeoportalClient``.
|
||||||
|
bbox: (minlon, minlat, maxlon, maxlat) EPSG:4326 — bounds для WFS-перечисления.
|
||||||
|
limit: максимум УНИКАЛЬНЫХ незакэшированных зон к фетчу (None = все). Для теста/частичного
|
||||||
|
прогона.
|
||||||
|
rate_delay_s: пауза (сек) между live-fetch'ами зон — вежливость к геопорталу.
|
||||||
|
city: ключ кэша (по умолчанию 'ekb').
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
{'enumerated_features', 'unique_zones', 'cached_before', 'fetched', 'upserted', 'failed'}.
|
||||||
|
"""
|
||||||
|
client = client or EKBGeoportalClient()
|
||||||
|
features = client.features_in_bbox("territorial_zone", bbox)
|
||||||
|
by_index = _dedup_zones_by_index(features)
|
||||||
|
logger.info(
|
||||||
|
"backfill_ekb_zone_regulations: фич=%d уникальных_зон=%d",
|
||||||
|
len(features),
|
||||||
|
len(by_index),
|
||||||
|
)
|
||||||
|
|
||||||
|
cached_before = 0
|
||||||
|
fetched = 0
|
||||||
|
upserted = 0
|
||||||
|
failed = 0
|
||||||
|
fetched_this_run = 0 # сколько зон уже сходили в live-fetch (для rate_delay между ними)
|
||||||
|
|
||||||
|
for zone_index, feat in by_index.items():
|
||||||
|
# Идемпотентность: уже в кэше → не дёргаем геопортал.
|
||||||
|
if get_cached_zone_regulation(db, zone_index, city=city) is not None:
|
||||||
|
cached_before += 1
|
||||||
|
continue
|
||||||
|
if limit is not None and fetched >= limit:
|
||||||
|
break
|
||||||
|
|
||||||
|
try:
|
||||||
|
point = _representative_point(feat.geometry)
|
||||||
|
if point is None:
|
||||||
|
continue
|
||||||
|
# Вежливость к геопорталу: пауза ПЕРЕД каждым fetch кроме первого.
|
||||||
|
if fetched_this_run > 0 and rate_delay_s > 0:
|
||||||
|
time.sleep(rate_delay_s)
|
||||||
|
fetched_this_run += 1
|
||||||
|
lon, lat = point
|
||||||
|
reg = client.zone_regulation_at(lon, lat)
|
||||||
|
fetched += 1
|
||||||
|
if reg is None or not reg.zone_index:
|
||||||
|
continue
|
||||||
|
if upsert_zone_regulation(db, reg, city=city) is not None:
|
||||||
|
db.commit()
|
||||||
|
upserted += 1
|
||||||
|
except Exception as exc:
|
||||||
|
# Per-zone изоляция: один сбойный фетч/апсёрт не должен ронять весь прогон.
|
||||||
|
logger.warning(
|
||||||
|
"backfill_ekb_zone_regulations: зона idx=%s failed: %s", zone_index, exc
|
||||||
|
)
|
||||||
|
failed += 1
|
||||||
|
continue
|
||||||
|
|
||||||
|
result = {
|
||||||
|
"enumerated_features": len(features),
|
||||||
|
"unique_zones": len(by_index),
|
||||||
|
"cached_before": cached_before,
|
||||||
|
"fetched": fetched,
|
||||||
|
"upserted": upserted,
|
||||||
|
"failed": failed,
|
||||||
|
}
|
||||||
|
logger.info("backfill_ekb_zone_regulations: %s", result)
|
||||||
|
return result
|
||||||
|
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
|
"backfill_ekb_zone_regulations",
|
||||||
"get_cached_zone_regulation",
|
"get_cached_zone_regulation",
|
||||||
"get_or_fetch_zone_regulation",
|
"get_or_fetch_zone_regulation",
|
||||||
"parse_limit_params",
|
"parse_limit_params",
|
||||||
|
|
|
||||||
|
|
@ -73,6 +73,7 @@ celery_app = Celery(
|
||||||
"app.workers.tasks.opportunity_harvest",
|
"app.workers.tasks.opportunity_harvest",
|
||||||
"app.workers.tasks.planning_harvest",
|
"app.workers.tasks.planning_harvest",
|
||||||
"app.workers.tasks.zone_regulation_refresh",
|
"app.workers.tasks.zone_regulation_refresh",
|
||||||
|
"app.workers.tasks.backfill_zone_regulations",
|
||||||
"app.workers.tasks.reservation_ingest",
|
"app.workers.tasks.reservation_ingest",
|
||||||
"app.workers.tasks.genplan_zones_sync",
|
"app.workers.tasks.genplan_zones_sync",
|
||||||
"app.workers.tasks.ekb_ppt_tep_sync",
|
"app.workers.tasks.ekb_ppt_tep_sync",
|
||||||
|
|
|
||||||
61
backend/app/workers/tasks/backfill_zone_regulations.py
Normal file
61
backend/app/workers/tasks/backfill_zone_regulations.py
Normal file
|
|
@ -0,0 +1,61 @@
|
||||||
|
"""Celery task: backfill кэша ПЗЗ-регламента ВСЕХ терзон ЕКБ (эпик #1881).
|
||||||
|
|
||||||
|
В отличие от ``zone_regulation_refresh.refresh_zone_regulations`` (ленивый прогрев, beat
|
||||||
|
~ежемесячно) эта задача — явный одноразовый/ad-hoc backfill, чтобы ``financial_estimate`` (#1881)
|
||||||
|
фаерил для большинства участков: ``max_far`` резолвится из ``zone_regulation_cache`` без
|
||||||
|
live-urbanCard в hot-пути. На момент задачи в кэше всего ~33 зоны → задача докэширует остальные.
|
||||||
|
|
||||||
|
Идемпотентно: ``backfill_ekb_zone_regulations`` пропускает уже закэшированные зоны (cached skip),
|
||||||
|
так что повторный запуск дешёв (только WFS-перечисление). Per-zone try/except в сервисе —
|
||||||
|
сбой одной зоны не валит прогон. ``rate_delay_s`` — вежливость к геопорталу.
|
||||||
|
|
||||||
|
Mirror conventions (zone_regulation_refresh.py / backend.md): SessionLocal()+try/finally,
|
||||||
|
sync (не async), commit per-upsert внутри сервиса.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
|
||||||
|
from app.core.db import SessionLocal
|
||||||
|
from app.services.scrapers.ekb_geoportal_client import EKBGeoportalClient
|
||||||
|
from app.services.site_finder.zone_regulation import (
|
||||||
|
_EKB_CITY_BBOX,
|
||||||
|
backfill_ekb_zone_regulations,
|
||||||
|
)
|
||||||
|
from app.workers.celery_app import celery_app
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
@celery_app.task(name="tasks.backfill_zone_regulations.backfill_zone_regulations")
|
||||||
|
def backfill_zone_regulations(
|
||||||
|
bbox: tuple[float, float, float, float] | None = None,
|
||||||
|
limit: int | None = None,
|
||||||
|
rate_delay_s: float = 1.0,
|
||||||
|
) -> dict[str, int]:
|
||||||
|
"""Backfill ``zone_regulation_cache`` по ВСЕМ терзонам ЕКБ (эпик #1881).
|
||||||
|
|
||||||
|
Args:
|
||||||
|
bbox: (minlon, minlat, maxlon, maxlat) EPSG:4326. None → дефолтный bbox ЕКБ.
|
||||||
|
limit: максимум уникальных незакэшированных зон к фетчу (None = все). Для smoke-теста.
|
||||||
|
rate_delay_s: пауза (сек) между live-fetch'ами зон — вежливость к геопорталу.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
{'enumerated_features', 'unique_zones', 'cached_before', 'fetched', 'upserted', 'failed'}.
|
||||||
|
"""
|
||||||
|
target_bbox = bbox if bbox is not None else _EKB_CITY_BBOX
|
||||||
|
client = EKBGeoportalClient()
|
||||||
|
db = SessionLocal()
|
||||||
|
try:
|
||||||
|
result = backfill_ekb_zone_regulations(
|
||||||
|
db,
|
||||||
|
client=client,
|
||||||
|
bbox=target_bbox,
|
||||||
|
limit=limit,
|
||||||
|
rate_delay_s=rate_delay_s,
|
||||||
|
)
|
||||||
|
logger.info("backfill_zone_regulations done: %s", result)
|
||||||
|
return result
|
||||||
|
finally:
|
||||||
|
db.close()
|
||||||
199
backend/tests/services/test_zone_regulation_backfill.py
Normal file
199
backend/tests/services/test_zone_regulation_backfill.py
Normal file
|
|
@ -0,0 +1,199 @@
|
||||||
|
"""Тесты backfill_ekb_zone_regulations (эпик #1881) — прогрев кэша регламента ВСЕХ терзон ЕКБ.
|
||||||
|
|
||||||
|
Сеть/БД не дёргаются: фейковый EKBGeoportalClient + фейковая Session + monkeypatch
|
||||||
|
``upsert_zone_regulation``/``get_cached_zone_regulation`` (логику апсёрта/чтения не тестируем тут —
|
||||||
|
она покрыта test_zone_regulation_extract). Проверяем:
|
||||||
|
- дедуп по urban_index (fetch только уникальных индексов),
|
||||||
|
- идемпотентность (cached skip → не фетчим),
|
||||||
|
- per-zone error isolation (одна зона падает → failed++, batch продолжается),
|
||||||
|
- limit,
|
||||||
|
- представительную точку (внутри полигона).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from app.services.scrapers.ekb_geoportal_client import EKBFeature, ZoneRegulation
|
||||||
|
from app.services.site_finder import zone_regulation as zr
|
||||||
|
|
||||||
|
|
||||||
|
def _poly() -> dict[str, Any]:
|
||||||
|
return {"type": "Polygon", "coordinates": [[[60, 56], [60.1, 56], [60.1, 56.1], [60, 56]]]}
|
||||||
|
|
||||||
|
|
||||||
|
def _feat(idx: str, geometry: dict[str, Any] | None = None) -> EKBFeature:
|
||||||
|
return EKBFeature(
|
||||||
|
feature_id=idx,
|
||||||
|
geometry=geometry if geometry is not None else _poly(),
|
||||||
|
properties={"urban_index": idx},
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _reg(idx: str) -> ZoneRegulation:
|
||||||
|
return ZoneRegulation(
|
||||||
|
zone_index=idx,
|
||||||
|
zone_full_name=f"{idx} зона",
|
||||||
|
main_vri=[],
|
||||||
|
conditional_vri=[],
|
||||||
|
auxiliary_vri=[],
|
||||||
|
limit_params=[],
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeClient:
|
||||||
|
"""Фейковый WFS-клиент: отдаёт заданные фичи, фиксирует резолвы.
|
||||||
|
|
||||||
|
Все тестовые полигоны идентичны → ``zone_regulation_at`` не различает зоны по координатам;
|
||||||
|
нам это и не нужно — проверяем счётчики/дедуп/число резолвов, а не маппинг точка→зона.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, feats: list[EKBFeature]) -> None:
|
||||||
|
self._feats = feats
|
||||||
|
self.resolved: list[tuple[float, float]] = []
|
||||||
|
|
||||||
|
def features_in_bbox(self, layer: str, bbox: Any) -> list[EKBFeature]:
|
||||||
|
return self._feats
|
||||||
|
|
||||||
|
def zone_regulation_at(self, lon: float, lat: float) -> ZoneRegulation:
|
||||||
|
self.resolved.append((lon, lat))
|
||||||
|
return _reg(f"Z-{len(self.resolved)}")
|
||||||
|
|
||||||
|
|
||||||
|
class _DB:
|
||||||
|
"""Фейковая Session: считает commit'ы, ничего не пишет."""
|
||||||
|
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self.commits = 0
|
||||||
|
|
||||||
|
def commit(self) -> None:
|
||||||
|
self.commits += 1
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture
|
||||||
|
def _no_sleep(monkeypatch: Any) -> None:
|
||||||
|
"""Не спим в тестах (rate_delay)."""
|
||||||
|
monkeypatch.setattr(zr.time, "sleep", lambda *_a, **_kw: None)
|
||||||
|
|
||||||
|
|
||||||
|
def test_dedup_and_counts(monkeypatch: Any, _no_sleep: None) -> None:
|
||||||
|
"""3 фичи (2 уникальных индекса + 1 дубль) → fetch/upsert только 2 уникальных."""
|
||||||
|
feats = [_feat("Ц-1"), _feat("Ц-1"), _feat("Ж-2")]
|
||||||
|
client = _FakeClient(feats)
|
||||||
|
upserts: list[ZoneRegulation] = []
|
||||||
|
monkeypatch.setattr(zr, "get_cached_zone_regulation", lambda *a, **kw: None)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
zr, "upsert_zone_regulation", lambda db, reg, **kw: upserts.append(reg) or {}
|
||||||
|
)
|
||||||
|
|
||||||
|
db = _DB()
|
||||||
|
result = zr.backfill_ekb_zone_regulations(db, client=client, rate_delay_s=0.0) # type: ignore[arg-type]
|
||||||
|
|
||||||
|
assert result["enumerated_features"] == 3
|
||||||
|
assert result["unique_zones"] == 2 # дубль Ц-1 схлопнут
|
||||||
|
assert result["cached_before"] == 0
|
||||||
|
assert result["fetched"] == 2 # только уникальные индексы
|
||||||
|
assert result["upserted"] == 2
|
||||||
|
assert result["failed"] == 0
|
||||||
|
assert len(client.resolved) == 2
|
||||||
|
assert db.commits == 2 # commit per-upsert
|
||||||
|
|
||||||
|
|
||||||
|
def test_idempotent_cached_skip(monkeypatch: Any, _no_sleep: None) -> None:
|
||||||
|
"""Если зона уже в кэше — НЕ фетчим (идемпотентность)."""
|
||||||
|
feats = [_feat("Ц-1"), _feat("Ж-2")]
|
||||||
|
client = _FakeClient(feats)
|
||||||
|
# Ц-1 закэширована, Ж-2 нет.
|
||||||
|
monkeypatch.setattr(
|
||||||
|
zr,
|
||||||
|
"get_cached_zone_regulation",
|
||||||
|
lambda db, zone_index, **kw: {"zone_index": zone_index} if zone_index == "Ц-1" else None,
|
||||||
|
)
|
||||||
|
upserts: list[ZoneRegulation] = []
|
||||||
|
monkeypatch.setattr(
|
||||||
|
zr, "upsert_zone_regulation", lambda db, reg, **kw: upserts.append(reg) or {}
|
||||||
|
)
|
||||||
|
|
||||||
|
result = zr.backfill_ekb_zone_regulations(_DB(), client=client, rate_delay_s=0.0) # type: ignore[arg-type]
|
||||||
|
|
||||||
|
assert result["cached_before"] == 1 # Ц-1 пропущена
|
||||||
|
assert result["fetched"] == 1 # только Ж-2
|
||||||
|
assert result["upserted"] == 1
|
||||||
|
assert len(client.resolved) == 1 # Ц-1 не дёргала геопортал
|
||||||
|
|
||||||
|
|
||||||
|
def test_per_zone_error_isolation(monkeypatch: Any, _no_sleep: None) -> None:
|
||||||
|
"""Сбой fetch одной зоны → failed++, batch продолжается для остальных."""
|
||||||
|
# Клиент кидает на ПЕРВОЙ резолвленной зоне (по порядку), вторую отдаёт нормально.
|
||||||
|
feats = [_feat("Ц-1"), _feat("Ж-2")]
|
||||||
|
|
||||||
|
calls = {"n": 0}
|
||||||
|
|
||||||
|
class _RaisingClient(_FakeClient):
|
||||||
|
def zone_regulation_at(self, lon: float, lat: float) -> ZoneRegulation:
|
||||||
|
self.resolved.append((lon, lat))
|
||||||
|
calls["n"] += 1
|
||||||
|
if calls["n"] == 1:
|
||||||
|
raise RuntimeError("geoportal 500 on first zone")
|
||||||
|
return _reg("Ж-2")
|
||||||
|
|
||||||
|
client = _RaisingClient(feats)
|
||||||
|
monkeypatch.setattr(zr, "get_cached_zone_regulation", lambda *a, **kw: None)
|
||||||
|
upserts: list[ZoneRegulation] = []
|
||||||
|
monkeypatch.setattr(
|
||||||
|
zr, "upsert_zone_regulation", lambda db, reg, **kw: upserts.append(reg) or {}
|
||||||
|
)
|
||||||
|
|
||||||
|
result = zr.backfill_ekb_zone_regulations(_DB(), client=client, rate_delay_s=0.0) # type: ignore[arg-type]
|
||||||
|
|
||||||
|
assert result["unique_zones"] == 2
|
||||||
|
assert result["failed"] == 1 # первая зона упала
|
||||||
|
assert result["upserted"] == 1 # вторая прошла → batch не остановился
|
||||||
|
assert len(client.resolved) == 2 # дошли до обеих зон
|
||||||
|
|
||||||
|
|
||||||
|
def test_limit_caps_fetches(monkeypatch: Any, _no_sleep: None) -> None:
|
||||||
|
"""limit ограничивает число live-fetch'ей незакэшированных зон."""
|
||||||
|
feats = [_feat("Ц-1"), _feat("Ж-2"), _feat("Ж-3")]
|
||||||
|
client = _FakeClient(feats)
|
||||||
|
monkeypatch.setattr(zr, "get_cached_zone_regulation", lambda *a, **kw: None)
|
||||||
|
upserts: list[ZoneRegulation] = []
|
||||||
|
monkeypatch.setattr(
|
||||||
|
zr, "upsert_zone_regulation", lambda db, reg, **kw: upserts.append(reg) or {}
|
||||||
|
)
|
||||||
|
|
||||||
|
result = zr.backfill_ekb_zone_regulations(_DB(), client=client, limit=2, rate_delay_s=0.0) # type: ignore[arg-type]
|
||||||
|
|
||||||
|
assert result["unique_zones"] == 3
|
||||||
|
assert result["fetched"] == 2 # обрезано лимитом
|
||||||
|
assert len(client.resolved) == 2
|
||||||
|
|
||||||
|
|
||||||
|
def test_representative_point_inside_polygon() -> None:
|
||||||
|
pt = zr._representative_point(_poly())
|
||||||
|
assert pt is not None
|
||||||
|
lon, lat = pt
|
||||||
|
assert 60.0 <= lon <= 60.1
|
||||||
|
assert 56.0 <= lat <= 56.1
|
||||||
|
|
||||||
|
|
||||||
|
def test_representative_point_none_on_bad_geom() -> None:
|
||||||
|
assert zr._representative_point(None) is None
|
||||||
|
assert zr._representative_point({"type": "Polygon", "coordinates": "garbage"}) is None
|
||||||
|
|
||||||
|
|
||||||
|
def test_dedup_prefers_feature_with_geometry() -> None:
|
||||||
|
"""Из двух фич одного индекса выбираем ту, у которой есть геометрия."""
|
||||||
|
feats = [_feat("Ц-1", geometry=None), _feat("Ц-1", geometry=_poly())]
|
||||||
|
by_index = zr._dedup_zones_by_index(feats)
|
||||||
|
assert set(by_index) == {"Ц-1"}
|
||||||
|
assert by_index["Ц-1"].geometry is not None
|
||||||
|
|
||||||
|
|
||||||
|
def test_blank_index_dropped() -> None:
|
||||||
|
"""Фичи с пустым/отсутствующим urban_index не попадают в дедуп."""
|
||||||
|
feats = [_feat("", geometry=_poly()), _feat("Ц-1")]
|
||||||
|
by_index = zr._dedup_zones_by_index(feats)
|
||||||
|
assert set(by_index) == {"Ц-1"}
|
||||||
49
backend/tests/workers/test_backfill_zone_regulations.py
Normal file
49
backend/tests/workers/test_backfill_zone_regulations.py
Normal file
|
|
@ -0,0 +1,49 @@
|
||||||
|
"""Тест Celery-таски backfill_zone_regulations (эпик #1881).
|
||||||
|
|
||||||
|
Сеть/БД не дёргаются: мокаем SessionLocal + backfill_ekb_zone_regulations (логика backfill
|
||||||
|
покрыта tests/services/test_zone_regulation_backfill). Проверяем что таска открывает/закрывает
|
||||||
|
сессию, прокидывает limit/rate_delay_s и возвращает результат сервиса.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
from app.workers.tasks import backfill_zone_regulations as task_mod
|
||||||
|
|
||||||
|
|
||||||
|
class _DB:
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self.closed = False
|
||||||
|
|
||||||
|
def close(self) -> None:
|
||||||
|
self.closed = True
|
||||||
|
|
||||||
|
|
||||||
|
def test_task_wires_args_and_closes_session(monkeypatch: Any) -> None:
|
||||||
|
db = _DB()
|
||||||
|
captured: dict[str, Any] = {}
|
||||||
|
|
||||||
|
def _fake_backfill(session: Any, **kwargs: Any) -> dict[str, int]:
|
||||||
|
captured["session"] = session
|
||||||
|
captured.update(kwargs)
|
||||||
|
return {
|
||||||
|
"enumerated_features": 5,
|
||||||
|
"unique_zones": 3,
|
||||||
|
"cached_before": 1,
|
||||||
|
"fetched": 2,
|
||||||
|
"upserted": 2,
|
||||||
|
"failed": 0,
|
||||||
|
}
|
||||||
|
|
||||||
|
monkeypatch.setattr(task_mod, "SessionLocal", lambda: db)
|
||||||
|
monkeypatch.setattr(task_mod, "EKBGeoportalClient", lambda *a, **kw: object())
|
||||||
|
monkeypatch.setattr(task_mod, "backfill_ekb_zone_regulations", _fake_backfill)
|
||||||
|
|
||||||
|
result = task_mod.backfill_zone_regulations(limit=2, rate_delay_s=0.0)
|
||||||
|
|
||||||
|
assert result["upserted"] == 2
|
||||||
|
assert captured["limit"] == 2
|
||||||
|
assert captured["rate_delay_s"] == 0.0
|
||||||
|
assert captured["session"] is db
|
||||||
|
assert db.closed is True # finally: db.close()
|
||||||
Loading…
Add table
Reference in a new issue