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