gendesign/backend/app/services/site_finder/quarter_dump_lookup.py
lekss361 7ca4d35e88 fix(site-finder): rename _make_empty_result → make_empty_result (public) per PR #116 review
M1 residual fix: parcels.py exception path использовал EMPTY_DUMP_RESULT
singleton ref вместо factory. Сейчас readonly access, но нарушает
documented invariant модуля.

Rename `_make_empty_result` → `make_empty_result` (public API), import в
parcels.py, использовать в try/except fallback. Каждый request получает
независимый dict — никаких shared references.

M4 (Redis SETNX dedup) + M5 (test fragility) — deferred per review,
documented в code/issue. Acceptable trade-offs:
- M4: UPSERT idempotency делает данные safe; burst-duplicate task'и тратят
  WAF traffic впустую но не повреждают данные.
- M5: docstring contractually describes positional-call order.

17/17 tests pass. ruff/format clean.

Per auto-review on aef8308.
2026-05-13 09:11:41 +03:00

417 lines
16 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.

"""Quarter dump lookup helper for analyze_parcel.
Sprint 1.1 item #4 (feat/analyze-uses-quarter-dump):
Читает nspd_quarter_dumps кеш и извлекает:
- nspd_zoning — территориальная зона ПЗЗ (G1) по centroid участка
- nspd_zouit_overlaps — список ЗОУИТ (G3) которые пересекаются с участком
- nspd_engineering_nearby — инженерные сооружения в 200м (I3)
- nspd_dump — freshness metadata (доступность, возраст, trigger флаг)
Если дамп отсутствует или устарел (>180 дней) — fire-and-forget
harvest_quarter.apply_async() и продолжает без dump-derived полей.
"""
from __future__ import annotations
import logging
from datetime import UTC, datetime, timedelta
from typing import Any
from sqlalchemy import text
from sqlalchemy.orm import Session
logger = logging.getLogger(__name__)
# Порог свежести дампа — 180 дней. Совпадает с beat-расписанием (раз в квартал
# на практике, но с запасом для редко запрашиваемых кварталов).
_DUMP_MAX_AGE_DAYS = 180
# Радиус поиска инженерных сооружений (метры) — договорённость #44 I3.
_ENGINEERING_RADIUS_M = 200
# Sentinel для isinstance-проверок и read-only fallback в parcels.py try/except.
# НИКОГДА не мутировать — использовать make_empty_result() для новых dict.
EMPTY_DUMP_RESULT: dict[str, Any] = {
"nspd_zoning": None,
"nspd_zouit_overlaps": [],
"nspd_engineering_nearby": [],
"nspd_dump": {
"available": False,
"fetched_at_utc": None,
"stale": False,
"harvest_triggered": False,
"total_features": None,
},
}
def make_empty_result(
*,
fetched_at_utc: str | None = None,
stale: bool = False,
harvest_triggered: bool = False,
total_features: int | None = None,
) -> dict[str, Any]:
"""Создаёт свежую копию empty-dump result с возможностью переопределить поля.
Вызывать вместо EMPTY_DUMP_RESULT напрямую — чтобы не мутировать singleton.
"""
return {
"nspd_zoning": None,
"nspd_zouit_overlaps": [],
"nspd_engineering_nearby": [],
"nspd_dump": {
"available": False,
"fetched_at_utc": fetched_at_utc,
"stale": stale,
"harvest_triggered": harvest_triggered,
"total_features": total_features,
},
}
def derive_quarter_cad(cad_num: str) -> str | None:
"""3-сегментный кадастровый квартал из любого кадастрового номера.
- 3-сегмент (квартал) '66:41:0204016''66:41:0204016'
- 4-сегмент (участок) '66:41:0204016:10''66:41:0204016'
- 5-сегмент (здание) '66:41:0204016:10:1''66:41:0204016'
- невалидный формат → None
"""
parts = cad_num.strip().split(":")
if len(parts) < 3:
return None
quarter = ":".join(parts[:3])
# Каждый сегмент — только цифры
if not all(p.isdigit() and len(p) >= 1 for p in parts[:3]):
return None
return quarter
def get_quarter_dump_data(
db: Session,
cad_num: str,
parcel_wkt: str | None,
) -> dict[str, Any]:
"""Читает quarter dump для квартала cad_num и возвращает НСПД-контекст.
Returns dict с ключами:
- nspd_zoning: dict | None — зона ПЗЗ из territorial_zones
- nspd_zouit_overlaps: list[dict] — ЗОУИТ пересечения
- nspd_engineering_nearby: list[dict] — инженерные сооружения в 200м
- nspd_dump: dict — freshness metadata
Если дамп отсутствует или устарел — вызывает harvest_quarter.apply_async()
(non-blocking) и возвращает пустые spatial поля + nspd_dump.available=False.
Если parcel_wkt=None — возвращает только freshness metadata (нет geom для
spatial queries).
"""
quarter = derive_quarter_cad(cad_num)
if quarter is None:
logger.warning("quarter_dump_lookup: cannot derive quarter from cad=%s", cad_num)
return make_empty_result()
# Читаем строку дампа из БД. Денормализованные счётчики слоёв используются
# для early-exit в spatial helpers (M2 mitigation).
row = db.execute(
text(
"""
SELECT quarter_cad,
fetched_at_utc,
total_features,
harvest_error,
territorial_zones_count,
zouit_count,
engineering_count
FROM nspd_quarter_dumps
WHERE quarter_cad = :q
"""
),
{"q": quarter},
).first()
now = datetime.now(UTC)
max_age = timedelta(days=_DUMP_MAX_AGE_DAYS)
if row is None:
# Дампа нет — ставим harvest в очередь
harvest_triggered = _trigger_harvest(quarter)
return make_empty_result(harvest_triggered=harvest_triggered)
fetched_at: datetime = row[1]
# Убедимся что timezone-aware для корректного сравнения
if fetched_at.tzinfo is None:
fetched_at = fetched_at.replace(tzinfo=UTC)
total_features: int | None = row[2]
harvest_error: str | None = row[3]
territorial_zones_count: int = row[4] or 0
zouit_count: int = row[5] or 0
engineering_count: int = row[6] or 0
is_stale = (now - fetched_at) > max_age
has_error = harvest_error is not None
# Устаревший или с ошибкой — триггерим повторный harvest
if is_stale or has_error:
harvest_triggered = _trigger_harvest(quarter)
return make_empty_result(
fetched_at_utc=fetched_at.isoformat(),
stale=is_stale,
harvest_triggered=harvest_triggered,
total_features=total_features,
)
# Свежий дамп без ошибок — извлекаем spatial данные
dump_meta: dict[str, Any] = {
"available": True,
"fetched_at_utc": fetched_at.isoformat(),
"stale": False,
"harvest_triggered": False,
"total_features": total_features,
}
if parcel_wkt is None:
# Нет геометрии участка — возвращаем только метаданные
return {
"nspd_zoning": None,
"nspd_zouit_overlaps": [],
"nspd_engineering_nearby": [],
"nspd_dump": dump_meta,
}
layer_counts = {
"territorial_zones_count": territorial_zones_count,
"zouit_count": zouit_count,
"engineering_count": engineering_count,
}
nspd_zoning = _get_zoning(db, quarter, parcel_wkt, layer_counts)
nspd_zouit = _get_zouit_overlaps(db, quarter, parcel_wkt, layer_counts)
nspd_engineering = _get_engineering_nearby(db, quarter, parcel_wkt, layer_counts)
return {
"nspd_zoning": nspd_zoning,
"nspd_zouit_overlaps": nspd_zouit,
"nspd_engineering_nearby": nspd_engineering,
"nspd_dump": dump_meta,
}
# ── Spatial helpers ───────────────────────────────────────────────────────────
def _get_zoning(
db: Session,
quarter: str,
parcel_wkt: str,
layer_counts: dict[str, int] | None = None,
) -> dict[str, Any] | None:
"""G1: ПЗЗ территориальная зона по centroid участка из dump.
Geometry в features_json в EPSG:3857 — ST_Transform на чтении → 4326.
layer_counts — денормализованные счётчики из строки дампа. Если
territorial_zones_count == 0 — пропускаем heavy jsonb_array_elements scan.
"""
if layer_counts is not None and layer_counts.get("territorial_zones_count", 1) == 0:
return None
try:
row = db.execute(
text(
"""
SELECT feat.value->'properties' AS zone_props
FROM nspd_quarter_dumps d,
jsonb_array_elements(d.features_json) AS feat(value)
WHERE d.quarter_cad = :q
AND feat.value->>'layer' = 'territorial_zones'
AND (feat.value->'geometry') IS NOT NULL
AND feat.value->>'geometry' != 'null'
AND ST_Intersects(
ST_Transform(
ST_SetSRID(
ST_GeomFromGeoJSON(feat.value->>'geometry'),
3857
),
4326
),
ST_Centroid(ST_GeomFromText(:wkt, 4326))
)
LIMIT 1
"""
),
{"q": quarter, "wkt": parcel_wkt},
).first()
if row is None:
return None
props: dict[str, Any] = row[0] if isinstance(row[0], dict) else {}
zone_code = props.get("reg_numb_border") or props.get("zone_code") or props.get("name")
zone_name = props.get("type_zone") or props.get("zone_name") or props.get("name")
return {
"zone_code": zone_code,
"zone_name": zone_name,
"source": "nspd-quarter-dump",
"raw_props": props,
}
except Exception as e:
logger.warning("nspd zoning query failed for quarter=%s: %s", quarter, e)
return None
def _get_zouit_overlaps(
db: Session,
quarter: str,
parcel_wkt: str,
layer_counts: dict[str, int] | None = None,
) -> list[dict[str, Any]]:
"""G3: ЗОУИТ которые пересекают участок.
Проверяем 5 групп: zouit_okn, zouit_engineering, zouit_natural,
zouit_protected, zouit_other.
layer_counts — денормализованные счётчики. Если zouit_count == 0 —
пропускаем heavy jsonb_array_elements scan.
"""
if layer_counts is not None and layer_counts.get("zouit_count", 1) == 0:
return []
try:
rows = db.execute(
text(
"""
SELECT feat.value->>'layer' AS layer,
feat.value->'properties' AS props
FROM nspd_quarter_dumps d,
jsonb_array_elements(d.features_json) AS feat(value)
WHERE d.quarter_cad = :q
AND feat.value->>'layer' LIKE 'zouit_%'
AND (feat.value->'geometry') IS NOT NULL
AND feat.value->>'geometry' != 'null'
AND ST_Intersects(
ST_Transform(
ST_SetSRID(
ST_GeomFromGeoJSON(feat.value->>'geometry'),
3857
),
4326
),
ST_GeomFromText(:wkt, 4326)
)
"""
),
{"q": quarter, "wkt": parcel_wkt},
).fetchall()
result: list[dict[str, Any]] = []
for r in rows:
layer: str = r[0] or ""
props: dict[str, Any] = r[1] if isinstance(r[1], dict) else {}
group_key = layer.removeprefix("zouit_")
result.append(
{
"group_key": group_key,
"layer": layer,
"subcategory": props.get("subcategory") or props.get("type_zone"),
"name": props.get("name") or props.get("object_name"),
"raw_props": props,
}
)
return result
except Exception as e:
logger.warning("nspd zouit query failed for quarter=%s: %s", quarter, e)
return []
def _get_engineering_nearby(
db: Session,
quarter: str,
parcel_wkt: str,
layer_counts: dict[str, int] | None = None,
) -> list[dict[str, Any]]:
"""I3: Инженерные сооружения в радиусе _ENGINEERING_RADIUS_M от centroid участка.
layer_counts — денормализованные счётчики. Если engineering_count == 0 —
пропускаем heavy jsonb_array_elements + ST_DWithin scan.
"""
if layer_counts is not None and layer_counts.get("engineering_count", 1) == 0:
return []
try:
rows = db.execute(
text(
"""
SELECT feat.value->'properties' AS props,
ST_Distance(
ST_Transform(
ST_SetSRID(
ST_GeomFromGeoJSON(feat.value->>'geometry'),
3857
),
4326
)::geography,
ST_Centroid(ST_GeomFromText(:wkt, 4326))::geography
) AS distance_m
FROM nspd_quarter_dumps d,
jsonb_array_elements(d.features_json) AS feat(value)
WHERE d.quarter_cad = :q
AND feat.value->>'layer' = 'engineering_structures'
AND (feat.value->'geometry') IS NOT NULL
AND feat.value->>'geometry' != 'null'
AND ST_DWithin(
ST_Transform(
ST_SetSRID(
ST_GeomFromGeoJSON(feat.value->>'geometry'),
3857
),
4326
)::geography,
ST_Centroid(ST_GeomFromText(:wkt, 4326))::geography,
:radius_m
)
ORDER BY distance_m ASC
LIMIT 20
"""
),
{"q": quarter, "wkt": parcel_wkt, "radius_m": _ENGINEERING_RADIUS_M},
).fetchall()
result: list[dict[str, Any]] = []
for r in rows:
props: dict[str, Any] = r[0] if isinstance(r[0], dict) else {}
distance_m = float(r[1]) if r[1] is not None else None
name = props.get("name") or props.get("object_name")
obj_type = props.get("object_type") or props.get("type_zone")
result.append(
{
"name": name,
"type": obj_type,
"distance_m": round(distance_m) if distance_m is not None else None,
"raw_props": props,
}
)
return result
except Exception as e:
logger.warning("nspd engineering query failed for quarter=%s: %s", quarter, e)
return []
# ── Harvest trigger ───────────────────────────────────────────────────────────
def _trigger_harvest(quarter: str) -> bool:
"""Fire-and-forget harvest_quarter.apply_async(). Возвращает True если enqueue OK.
Ленивый импорт чтобы избежать circular import (tasks → services → tasks).
Known limitation: burst из N concurrent analyze_parcel запросов на один
и тот же ещё не закешированный квартал может поставить N одинаковых задач
в очередь (нет дедупликации). UPSERT в harvest_quarter идемпотентен, поэтому
данные не портятся, но WAF traffic тратится впустую.
TODO: добавить Redis SETNX lock с TTL перед apply_async — отдельная issue/PR.
"""
try:
from app.workers.tasks.nspd_sync import harvest_quarter
harvest_quarter.apply_async(args=[quarter], kwargs={"region_code": 66})
logger.info("quarter dump harvest triggered for quarter=%s", quarter)
return True
except Exception as e:
logger.warning("failed to trigger harvest for quarter=%s: %s", quarter, e)
return False