* feat(site-finder): integrate nspd_quarter_dumps cache в analyze_parcel (#94 Sprint 1.1 #4 FINAL) Замыкает Sprint 1.1 из #94 part 2 plan. После этого PR пользователь видит свежие НСПД данные в UI (frontend integration — отдельный PR). Backend (new app/services/site_finder/quarter_dump_lookup.py): - `derive_quarter_cad(cad_num)` — 3/4/5-сегмент → quarter (3-segment) - `get_quarter_dump_data(db, cad_num, parcel_wkt)` — main entrypoint: - Reads nspd_quarter_dumps row для derived quarter - Freshness threshold: 180 days - Missing/stale/harvest_error → trigger harvest_quarter.apply_async() fire- and-forget (lazy import против circular), return EMPTY_DUMP_RESULT - Fresh + parcel_wkt=None → metadata only (no spatial queries) - Fresh + geometry → 3 spatial queries via jsonb_array_elements + ST_Transform (3857→4326) + ST_Intersects / ST_DWithin - 3 private helpers: - `_get_zoning` — point-in-polygon parcel centroid vs territorial_zones, LIMIT 1 - `_get_zouit_overlaps` — все zouit_% layers пересекающиеся с parcel - `_get_engineering_nearby` — engineering_structures в 200m, sorted by distance - `EMPTY_DUMP_RESULT` module-level constant — DRY для no-dump fallback (used in get_quarter_dump_data internal + analyze_parcel try/except wrap) Backend (parcels.py): - Import EMPTY_DUMP_RESULT + get_quarter_dump_data - Call wrapped в try/except — если nspd_quarter_dumps недоступна (DB timeout / table missing) → EMPTY_DUMP_RESULT fallback вместо 500 (consistent с resilience pattern других optional fetches) - Response gets 4 new fields: - nspd_zoning: dict | None (G1 ПЗЗ — zone_code, zone_name, source) - nspd_zouit_overlaps: list[dict] (G3 — overlaps в parcel, per ЗОУИТ group) - nspd_engineering_nearby: list[dict] (I3 — engineering structures в 200m) - nspd_dump: dict (freshness metadata — available, fetched_at_utc, stale, harvest_triggered, total_features) Tests: 13 new в test_quarter_dump_lookup.py (mock-based, no real DB): - derive_quarter_cad 5 edge cases (3seg, 4seg, 5seg, invalid, whitespace) - get_quarter_dump no_row → harvest triggered - stale (>180d) → harvest triggered, stale=True - harvest_error row → retry harvest triggered - parcel_wkt=None → metadata only (1 DB call) - fresh + zoning extraction - fresh + zouit_overlaps list - fresh integration: все 4 keys present 47 pre-existing tests still pass. Code review (code-reviewer pre-push): MINOR, 0 blocking. Applied 2 of 4: - ✅ #1: try/except wrap around get_quarter_dump_data в analyze_parcel (защита от DB unavailability) + DRY через EMPTY_DUMP_RESULT module const - ✅ #2: removed redundant nspd_zoning.fetched_at_utc (DRY — freshness в nspd_dump.fetched_at_utc) - ⏭ Deferred (acceptable): #3 ad-hoc harvest_quarter retry cooldown для harvest_error rows (только при high traffic + persistent NSPD errors); #4 raw_props в response — tech debt, убрать вместе с frontend PR Performance note: 3 spatial queries per analyze adds ~10-50ms on typical ~100-feature quarter. Mitigation if quarters grow dense: materialized per-layer sub-table (отдельная DB issue). Closes Sprint 1.1 part of #94. Frontend rendering этих 4 полей — отдельный PR (next: #112 / #115 / #114). * fix(site-finder): address PR #116 auto-review M1-M5 M1 (mutation risk): replace EMPTY_DUMP_RESULT direct refs with _make_empty_result() factory. dict(...) shallow copy left nested nspd_dump shared by reference across concurrent requests — single mutation pollutes module sentinel for all subsequent calls. Now каждый caller gets independent dict. M2 (O(N) spatial scan): SELECT extended denormalized counts (territorial_zones_count, zouit_count, engineering_count). Each spatial helper accepts layer_counts and early-returns when count=0 — skips heavy jsonb_array_elements + ST_Transform/ST_Intersects scan entirely. Critical для quarters с 2000+ features. M4 (documentation): _trigger_harvest docstring describes known burst/no-dedup limitation + TODO Redis SETNX (отдельный PR). M5 (test fragility): _make_db_mock_with_spatial docstring describes positional-call contract — db.execute order (0=dump, 1=zoning, 2=zouit, 3=engineering) и зависимость от count-values. +4 new tests (17 total, all pass): - test_make_empty_result_returns_independent_copies (mutation safety) - test_make_empty_result_overrides - test_early_exit_all_counts_zero_no_spatial_queries - test_early_exit_partial_counts Per auto-review on3068a9c. * 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 onaef8308. --------- Co-authored-by: lekss361 <claudestars@proton.me>
417 lines
16 KiB
Python
417 lines
16 KiB
Python
"""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
|