gendesign/backend/app/services/site_finder/rosseti_wfs_loader.py
Light1YT e285cbe33a
All checks were successful
CI / changes (pull_request) Successful in 8s
CI / frontend-tests (pull_request) Successful in 1m17s
CI / openapi-codegen-check (pull_request) Successful in 2m6s
CI / backend-tests (pull_request) Successful in 14m28s
feat(site-finder): резервы мощности для ТП — электро+вода на карте §3 (#2119 Фаза A)
Наполнение карты точками подключения с ХАРАКТЕРИСТИКАМИ (свободная мощность)
из верифицированных источников (research #2119). Все источники гео-блокируют
не-RU IP — лоадеры выполняются на проде (Celery weekly + ручной docker exec).

Backend:
- Миграция 180: power_supply_centers (ПС 35-220кВ, точки + резерв МВА),
  power_tp_rp_reserves (10.7k ТП/РП 6-10кВ, геокод — Фаза B),
  water_supply_reserves (ЦСВ/ЦСК Водоканала). UNIQUE NULLS NOT DISTINCT на
  nullable-ключах (грабли #140).
- rosseti_wfs_loader: открытый WFS портал-тп.рф (punycode) → 488 ПС Свердл обл,
  индекс загрузки 256184/6/8 → open/limited/closed, voltage из имени.
- rosseti_reserve_loader: раскрытие ПП№24 (xlsx) — ЦП 35-110кВ (строка 8+)
  матчится к WFS-точкам по нормализованному имени; ТП/РП <35кВ (строка 9+),
  кВА-санити (>2.5 МВА для ТП → ÷1000), «н/д»→NULL, SAVEPOINT per-row.
- vodokanal_reserve_loader: DOCX-раскрытие ПП№6 (stdlib parse, vMerge
  forward-fill), отрицательные резервы (город: дефицит) сохраняются со знаком.
- Endpoint GET /{cad}/connection-capacity: ПС в радиусе (ST_DWithin geography)
  + summary + вода за последний период per-system_kind. Celery: вт 04:00/04:30.

Frontend (§3):
- Карточка «Ресурсные резервы»: ближайший ЦП с резервом МВА + статус-бейдж;
  вода/канализация с красным «дефицит» при отрицательном резерве + примечание.
- Слой «Центры питания (резерв)» на карте: цвет по индексу загрузки
  (зелёный/янтарь/красный), попап с напряжением/мощностью/загрузкой/резервом,
  легенда под картой. Числа через toFiniteNumber (Decimal→str coercion).

api-types.ts перегенерён офлайн (+131/-0, аддитивно). 39 тестов зелёные.
Deep-review: ⚠️ MINOR → все 3 замечания исправлены (NULLS NOT DISTINCT,
per-kind MAX(period), settlement-комментарий).
2026-07-02 12:45:55 +05:00

269 lines
12 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Загрузчик центров питания (ЦП) Россети из WFS — свободная мощность §3 (#2119).
No-B2B верифицированный источник: WFS геосервера портала-тп.рф отдаёт ~488
центров питания (ПС 35/110 кВ) по Свердловской области с координатами (WGS84,
округлены ~100 м), классом напряжения и индексом загрузки (открыт/ограничен/
закрыт для техприсоединения). UPSERT-ит в ``power_supply_centers``.
Источник ГЕО-БЛОКИРУЕТ non-RU IP → загрузчик РАБОТАЕТ НА ПРОДЕ (Celery weekly /
manual docker exec). Punycode-host ОБЯЗАТЕЛЕН — кириллический алиас портал-тп.рф
не резолвится корректно из curl/httpx.
Структура зеркалит utility_infrastructure_loader: httpx с явным таймаутом,
per-row SAVEPOINT при UPSERT (битая фича не валит weekly-sync).
Также экспортирует ``normalize_sc_name`` — общий нормализатор имён ЦП, которым
пользуется и ``rosseti_reserve_loader`` (матч xlsx-резерва по имени), и endpoint.
"""
import hashlib
import json
import logging
import re
import httpx
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.core.db import SessionLocal
logger = logging.getLogger(__name__)
# Punycode-host ОБЯЗАТЕЛЕН (кириллический алиас портал-тп.рф не для curl/httpx).
# xn----7sb7akeedqd.xn--p1ai == портал-тп.рф.
WFS_BASE_URL = "https://xn----7sb7akeedqd.xn--p1ai/geoserver/wfs"
# Слой центров питания (полный вид с резервами/напряжением).
_WFS_TYPE_NAME = "gisosslabs:sc_points_fullview"
# CQL-фильтр по региону — только Свердловская область.
_WFS_CQL_FILTER = "rg_code='SverdlovskOblast'"
# Таймаут WFS-запроса (сек). ~488 фич — с запасом.
_WFS_TIMEOUT = 60
# sc_indexload_id (справочник Россети) → индекс загрузки ЦП. Неизвестный id →
# None (сырое значение остаётся в raw). Верифицированные id из справочника:
_INDEXLOAD_MAP: dict[int, str] = {
256184: "open", # открыт для техприсоединения (есть резерв)
256186: "limited", # ограниченно
256188: "closed", # закрыт (резерва нет)
}
# «ПС»/«подстанция» + класс напряжения-префикс — срезаем при нормализации имени.
# Оба компонента опциональны, но хотя бы один должен присутствовать (иначе матч
# пустой и sub ничего не делает). Примеры:
# «ПС 110/35/10 Нижне-Исетская» → «нижне-исетская»
# «ПС Уктус» → «уктус» (только «ПС», без напряжения)
# «110/10 Уктус» → «уктус» (только напряжение)
_VOLTAGE_PREFIX_RE = re.compile(
r"^\s*(?:(?:пс|подстанци\w*)\s*)?(?:[\d/.,]+\s*(?:кв)?\s*)?",
re.IGNORECASE,
)
# Класс напряжения внутри имени ЦП, напр. «110/35/10» или «110/10».
_VOLTAGE_CLASS_RE = re.compile(r"\b(\d{1,3}(?:[/.]\d{1,3}){1,3})\b")
_MULTISPACE_RE = re.compile(r"\s+")
def normalize_sc_name(name: str | None) -> str:
"""Нормализует имя ЦП для матча xlsx-резерва с WFS-фичей.
Правила: lower, ё→е, срез «ПС»/подстанция + класс напряжения-префикса,
удаление кавычек, схлопывание тире/пробелов. Пусто → ''.
Examples:
«ПС 110/35/10 Нижне-Исетская» → «нижне-исетская»
«Нижне — Исетская» → «нижне-исетская»
«ПС "Уктус"» → «уктус»
"""
if not name:
return ""
s = name.strip().lower().replace("ё", "е")
s = s.replace("«", "").replace("»", "").replace('"', "").replace("'", "")
# Срез voltage/«ПС»-префикса в начале имени (не трогаем цифры внутри имени).
s = _VOLTAGE_PREFIX_RE.sub("", s, count=1)
# Схлопываем все виды тире (—, , -) в одиночный дефис, пробелы вокруг него убираем.
s = re.sub(r"\s*[—–-]\s*", "-", s)
s = _MULTISPACE_RE.sub(" ", s).strip()
return s
def parse_voltage_class(name: str | None) -> str | None:
"""Извлекает класс напряжения из текста имени ЦП («ПС 110/35/10 …» → '110/35/10').
Возвращает первую последовательность вида «NNN/NNN[/NNN…]». None — если нет.
"""
if not name:
return None
m = _VOLTAGE_CLASS_RE.search(name)
if not m:
return None
return m.group(1).replace(".", "/")
def _stable_external_id(feature: dict, props: dict) -> str:
"""Стабильный external_id фичи: feature['id'] или хэш ключевых полей.
WFS обычно отдаёт стабильный ``feature['id']``; если его нет — детерминированный
sha1 по (sc_name, координаты) чтобы UPSERT оставался идемпотентным.
"""
fid = feature.get("id")
if fid:
return str(fid)
geom = feature.get("geometry") or {}
coords = geom.get("coordinates")
seed = f"{props.get('sc_name', '')}|{coords}"
# sha1 здесь — стабильный дедуп-id фичи, не криптография.
return "h:" + hashlib.sha1(seed.encode("utf-8")).hexdigest()[:16]
def _map_load_index(props: dict) -> str | None:
"""sc_indexload_id → 'open'|'limited'|'closed'|None (неизвестный → None)."""
raw = props.get("sc_indexload_id")
if raw is None:
return None
try:
return _INDEXLOAD_MAP.get(int(raw))
except (TypeError, ValueError):
return None
def fetch_power_supply_centers() -> list[dict]:
"""Тянет WFS-фичи центров питания Свердловской области (GeoJSON Feature-list).
RUN-ON-PROD: источник гео-блокирует non-RU IP. httpx с явным таймаутом,
verify по умолчанию (валидный TLS у геосервера). Возвращает list feature-dict'ов
``{"id", "geometry": {...}, "properties": {...}}``.
"""
params = {
"service": "WFS",
"version": "1.1.0",
"request": "GetFeature",
"typeName": _WFS_TYPE_NAME,
"outputFormat": "json",
"CQL_FILTER": _WFS_CQL_FILTER,
}
resp = httpx.get(WFS_BASE_URL, params=params, timeout=_WFS_TIMEOUT)
resp.raise_for_status()
data = resp.json()
features: list[dict] = (data or {}).get("features") or []
logger.info("rosseti_wfs: загружено центров питания=%d", len(features))
return features
def _point_geom_sql(feature: dict) -> tuple[str, dict] | None:
"""SQL-выражение + params для geom из GeoJSON Point. None — если не Point/нет коорд."""
geom = feature.get("geometry") or {}
if geom.get("type") != "Point":
return None
coords = geom.get("coordinates")
if not coords or len(coords) < 2:
return None
lon, lat = coords[0], coords[1]
if lon is None or lat is None:
return None
return "ST_SetSRID(ST_MakePoint(:lon, :lat), 4326)", {"lon": lon, "lat": lat}
def load_power_supply_centers(db: Session | None = None) -> dict[str, int]:
"""Тянет WFS + UPSERT-ит центры питания в power_supply_centers.
UPSERT по UNIQUE(source, external_id). Per-row SAVEPOINT — битая фича не валит
weekly-sync. ``db`` принимается для совместимости сигнатуры; если None — своя
SessionLocal (WFS-fetch + запись в БД sync, как utility_infrastructure_loader).
Returns: счётчики fetched/inserted/updated/skipped.
"""
features = fetch_power_supply_centers()
owns_session = db is None
if db is None:
db = SessionLocal()
inserted = 0
updated = 0
skipped = 0
fetched = len(features)
try:
for feature in features:
props: dict = feature.get("properties") or {}
sc_name = props.get("sc_name")
if not sc_name:
skipped += 1
continue
geom_pair = _point_geom_sql(feature)
geom_sql = geom_pair[0] if geom_pair else "NULL"
geom_params = geom_pair[1] if geom_pair else {}
params: dict = {
"source": "rosseti_wfs",
"external_id": _stable_external_id(feature, props),
"sc_name": sc_name,
"sc_name_norm": normalize_sc_name(sc_name),
"dzo_name": props.get("dzo_name"),
"org_name": props.get("org_name"),
"branch_name": props.get("br_name"),
"voltage_class": parse_voltage_class(sc_name),
"load_index": _map_load_index(props),
"municipality": props.get("municipality"),
"place_city": props.get("sc_placecity"),
"raw": json.dumps(feature, ensure_ascii=False),
**geom_params,
}
try:
with db.begin_nested(): # SAVEPOINT — откат только этой фичи
result = db.execute(
text(f"""
INSERT INTO power_supply_centers
(source, external_id, sc_name, sc_name_norm, dzo_name,
org_name, branch_name, voltage_class, load_index,
municipality, place_city, geom, raw, fetched_at)
VALUES (
:source, :external_id, :sc_name, :sc_name_norm, :dzo_name,
:org_name, :branch_name, :voltage_class, :load_index,
:municipality, :place_city, {geom_sql},
CAST(:raw AS jsonb), NOW()
)
ON CONFLICT (source, external_id) DO UPDATE
SET sc_name = EXCLUDED.sc_name,
sc_name_norm = EXCLUDED.sc_name_norm,
dzo_name = EXCLUDED.dzo_name,
org_name = EXCLUDED.org_name,
branch_name = EXCLUDED.branch_name,
voltage_class = EXCLUDED.voltage_class,
load_index = EXCLUDED.load_index,
municipality = EXCLUDED.municipality,
place_city = EXCLUDED.place_city,
geom = EXCLUDED.geom,
raw = EXCLUDED.raw,
fetched_at = NOW()
RETURNING (xmax = 0) AS is_insert
"""),
params,
).scalar()
if result:
inserted += 1
else:
updated += 1
except Exception as e:
logger.warning("rosseti_wfs upsert failed for %r: %s", sc_name, e)
skipped += 1
db.commit()
except Exception as e:
db.rollback()
logger.exception("rosseti_wfs: unexpected error, outer tx rolled back: %s", e)
raise
finally:
if owns_session:
db.close()
result_dict = {
"fetched": fetched,
"inserted": inserted,
"updated": updated,
"skipped": skipped,
}
logger.info("rosseti_wfs load done: %s", result_dict)
return result_dict