8 changed files with 807 additions and 40 deletions
|
|
@ -104,6 +104,21 @@ class NspdBulkRateLimitError(NspdBulkError):
|
||||||
"""HTTP 429 Rate limit — исчерпаны retries после backoff."""
|
"""HTTP 429 Rate limit — исчерпаны retries после backoff."""
|
||||||
|
|
||||||
|
|
||||||
|
class NspdBulkServerError(NspdBulkError):
|
||||||
|
"""HTTP 5xx / WMS ServiceException — server-side ошибка NSPD.
|
||||||
|
|
||||||
|
Issue #252: NSPD WMS отдаёт 500 + тело с <ServiceException> на десятки cells
|
||||||
|
per quarter (повторяющиеся в GlitchTip BACKEND-8/10/11/12/14/16/17/24/36/37).
|
||||||
|
Это transient server-side noise, не наш баг и не клиентская ошибка (4xx).
|
||||||
|
Выделяем в отдельный подкласс, чтобы:
|
||||||
|
- caller (bulk_harvest_quarter_task) ретраил весь квартал — autoretry_for
|
||||||
|
ловит подкласс NspdBulkError, dont_autoretry_for=(NspdBulkWafError,) его
|
||||||
|
не исключает;
|
||||||
|
- grid-walk считал per-layer fail-rate и поднимал layer_X_failed-флаг,
|
||||||
|
если ВСЕ cells слоя упали с 500 (см. _grid_walk_category).
|
||||||
|
"""
|
||||||
|
|
||||||
|
|
||||||
# ── NSPDBulkClient ─────────────────────────────────────────────────────────────
|
# ── NSPDBulkClient ─────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -169,7 +184,8 @@ class NSPDBulkClient:
|
||||||
Raises:
|
Raises:
|
||||||
NspdBulkWafError при 403.
|
NspdBulkWafError при 403.
|
||||||
NspdBulkRateLimitError при 429 после исчерпания retries.
|
NspdBulkRateLimitError при 429 после исчерпания retries.
|
||||||
NspdBulkError при прочих 4xx/5xx.
|
NspdBulkServerError при 5xx или WMS <ServiceException> в теле.
|
||||||
|
NspdBulkError при прочих 4xx.
|
||||||
"""
|
"""
|
||||||
if self._client is None or self._sem is None:
|
if self._client is None or self._sem is None:
|
||||||
raise RuntimeError("NSPDBulkClient не инициализирован — используй async with")
|
raise RuntimeError("NSPDBulkClient не инициализирован — используй async with")
|
||||||
|
|
@ -207,10 +223,29 @@ class NSPDBulkClient:
|
||||||
await asyncio.sleep(backoff)
|
await asyncio.sleep(backoff)
|
||||||
continue
|
continue
|
||||||
|
|
||||||
|
# Issue #252: 5xx — server-side ошибка NSPD (transient). Отдельный
|
||||||
|
# подкласс, чтобы caller ретраил квартал и считал per-layer fail-rate.
|
||||||
|
if resp.status_code >= 500:
|
||||||
|
body_preview = resp.text[:300]
|
||||||
|
raise NspdBulkServerError(
|
||||||
|
f"HTTP {resp.status_code} server error: {url} — {body_preview}"
|
||||||
|
)
|
||||||
|
|
||||||
if resp.status_code >= 400:
|
if resp.status_code >= 400:
|
||||||
body_preview = resp.text[:300]
|
body_preview = resp.text[:300]
|
||||||
raise NspdBulkError(f"HTTP {resp.status_code}: {url} — {body_preview}")
|
raise NspdBulkError(f"HTTP {resp.status_code}: {url} — {body_preview}")
|
||||||
|
|
||||||
|
# Issue #252: NSPD WMS (GeoServer) иногда отдаёт HTTP 200 с XML-телом
|
||||||
|
# <ServiceException> вместо запрошенного application/json (например при
|
||||||
|
# внутренней ошибке рендера слоя). resp.json() в этом случае кинул бы
|
||||||
|
# JSONDecodeError, маскируя server-side природу сбоя. Детектим маркер
|
||||||
|
# в теле и поднимаем NspdBulkServerError — единый retry/skip-путь с 5xx.
|
||||||
|
body = resp.text
|
||||||
|
if "ServiceException" in body or "<ServiceExceptionReport" in body:
|
||||||
|
raise NspdBulkServerError(
|
||||||
|
f"WMS ServiceException (HTTP {resp.status_code}): {url} — {body[:300]}"
|
||||||
|
)
|
||||||
|
|
||||||
return resp.json()
|
return resp.json()
|
||||||
|
|
||||||
# ── 1. search_by_quarter ──────────────────────────────────────────────────
|
# ── 1. search_by_quarter ──────────────────────────────────────────────────
|
||||||
|
|
@ -570,5 +605,6 @@ __all__ = [
|
||||||
"NSPDBulkClient",
|
"NSPDBulkClient",
|
||||||
"NspdBulkError",
|
"NspdBulkError",
|
||||||
"NspdBulkRateLimitError",
|
"NspdBulkRateLimitError",
|
||||||
|
"NspdBulkServerError",
|
||||||
"NspdBulkWafError",
|
"NspdBulkWafError",
|
||||||
]
|
]
|
||||||
|
|
|
||||||
|
|
@ -23,14 +23,14 @@ import hashlib
|
||||||
import json
|
import json
|
||||||
import logging
|
import logging
|
||||||
from collections.abc import Callable
|
from collections.abc import Callable
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass, field
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
from app.schemas.nspd_bulk import NSPDBulkFeature, QuarterSnapshot
|
from app.schemas.nspd_bulk import NSPDBulkFeature, QuarterSnapshot
|
||||||
from app.scrapers.nspd_bulk_client import NSPDBulkClient
|
from app.scrapers.nspd_bulk_client import NSPDBulkClient, NspdBulkServerError
|
||||||
from app.services.cadastre.grid_geometry import generate_grid_click_points, quarter_bbox_3857
|
from app.services.cadastre.grid_geometry import generate_grid_click_points, quarter_bbox_3857
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
@ -72,6 +72,10 @@ class HarvestResult:
|
||||||
zouit_upserted: int = 0
|
zouit_upserted: int = 0
|
||||||
snapshot_requests: int = 0 # search_by_quarter (Phase 1) + per-cat probes (Phase 1.5)
|
snapshot_requests: int = 0 # search_by_quarter (Phase 1) + per-cat probes (Phase 1.5)
|
||||||
grid_walk_requests: int = 0 # wms_feature_info (Phase 2-3)
|
grid_walk_requests: int = 0 # wms_feature_info (Phase 2-3)
|
||||||
|
# Issue #252: layer_id'ы, чей grid-walk целиком упал на server-side 500/
|
||||||
|
# ServiceException (все cells сбойнули). Квартал НЕ валится — слой
|
||||||
|
# пропускается, факт фиксируется в cadastre_jobs.phase_state.harvest_meta.
|
||||||
|
failed_layers: list[int] = field(default_factory=list)
|
||||||
phase_state: dict[str, Any] | None = None
|
phase_state: dict[str, Any] | None = None
|
||||||
|
|
||||||
@property
|
@property
|
||||||
|
|
@ -214,7 +218,7 @@ async def harvest_quarter(
|
||||||
cat_label = "parcels" if cat_id == CAT_PARCEL else "buildings"
|
cat_label = "parcels" if cat_id == CAT_PARCEL else "buildings"
|
||||||
update_progress({"phase": f"grid_walk_{cat_label}_started", "quarter": quarter})
|
update_progress({"phase": f"grid_walk_{cat_label}_started", "quarter": quarter})
|
||||||
|
|
||||||
discovered, n_requests = await _grid_walk_category(
|
discovered, n_requests, layer_failed = await _grid_walk_category(
|
||||||
db=db,
|
db=db,
|
||||||
client=client,
|
client=client,
|
||||||
quarter=quarter,
|
quarter=quarter,
|
||||||
|
|
@ -228,14 +232,20 @@ async def harvest_quarter(
|
||||||
else:
|
else:
|
||||||
result.buildings_upserted += discovered
|
result.buildings_upserted += discovered
|
||||||
|
|
||||||
|
# Issue #252: per-layer skip — слой целиком сбойнул на server-side 500/
|
||||||
|
# ServiceException. Фиксируем в harvest_meta (cadastre_jobs.phase_state),
|
||||||
|
# квартал продолжаем (остальные слои и фазы отрабатывают штатно).
|
||||||
|
done_progress: dict[str, Any] = {
|
||||||
|
"phase": f"grid_walk_{cat_label}_done",
|
||||||
|
"quarter": quarter,
|
||||||
|
f"{cat_label}_grid_discovered": discovered,
|
||||||
|
}
|
||||||
|
if layer_failed:
|
||||||
|
result.failed_layers.append(cat_id)
|
||||||
|
done_progress["harvest_meta"] = {f"layer_{cat_id}_failed": True}
|
||||||
|
|
||||||
db.commit()
|
db.commit()
|
||||||
update_progress(
|
update_progress(done_progress)
|
||||||
{
|
|
||||||
"phase": f"grid_walk_{cat_label}_done",
|
|
||||||
"quarter": quarter,
|
|
||||||
f"{cat_label}_grid_discovered": discovered,
|
|
||||||
}
|
|
||||||
)
|
|
||||||
|
|
||||||
# ── Phase 2.5: grid-walk для territorial_zones (ПЗЗ, layer 875838) ────────
|
# ── Phase 2.5: grid-walk для territorial_zones (ПЗЗ, layer 875838) ────────
|
||||||
# Выполняем после основного grid-walk (Phase 2-3). Требует bbox квартала.
|
# Выполняем после основного grid-walk (Phase 2-3). Требует bbox квартала.
|
||||||
|
|
@ -266,18 +276,26 @@ async def harvest_quarter(
|
||||||
logger.warning("harvest_quarter: geom auto-heal failed for %s: %s", quarter, e)
|
logger.warning("harvest_quarter: geom auto-heal failed for %s: %s", quarter, e)
|
||||||
db.commit()
|
db.commit()
|
||||||
|
|
||||||
update_progress({"phase": "done", "quarter": quarter})
|
# Issue #252: финальный phase_state несёт АГРЕГИРОВАННЫЙ harvest_meta по всем
|
||||||
result.phase_state = {"phase": "done", "quarter": quarter}
|
# сбойным слоям. progress_cb мержит phase_state через JSONB `||` (shallow) —
|
||||||
|
# per-layer done-апдейты перетёрли бы harvest_meta друг друга, поэтому в
|
||||||
|
# терминальном 'done' собираем полный список → итоговая строка корректна.
|
||||||
|
done_state: dict[str, Any] = {"phase": "done", "quarter": quarter}
|
||||||
|
if result.failed_layers:
|
||||||
|
done_state["harvest_meta"] = {f"layer_{lid}_failed": True for lid in result.failed_layers}
|
||||||
|
update_progress(done_state)
|
||||||
|
result.phase_state = done_state
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"harvest_quarter done: quarter=%s parcels=%d buildings=%d "
|
"harvest_quarter done: quarter=%s parcels=%d buildings=%d "
|
||||||
"constructions=%d zouit=%d grid_requests=%d",
|
"constructions=%d zouit=%d grid_requests=%d failed_layers=%s",
|
||||||
quarter,
|
quarter,
|
||||||
result.parcels_upserted,
|
result.parcels_upserted,
|
||||||
result.buildings_upserted,
|
result.buildings_upserted,
|
||||||
result.constructions_upserted,
|
result.constructions_upserted,
|
||||||
result.zouit_upserted,
|
result.zouit_upserted,
|
||||||
result.grid_walk_requests,
|
result.grid_walk_requests,
|
||||||
|
result.failed_layers or None,
|
||||||
)
|
)
|
||||||
return result
|
return result
|
||||||
|
|
||||||
|
|
@ -294,7 +312,7 @@ async def _grid_walk_category(
|
||||||
tile_size: int = 512,
|
tile_size: int = 512,
|
||||||
update_progress: Callable[[dict[str, Any]], None] | None = None,
|
update_progress: Callable[[dict[str, Any]], None] | None = None,
|
||||||
heartbeat_every: int = 50,
|
heartbeat_every: int = 50,
|
||||||
) -> tuple[int, int]:
|
) -> tuple[int, int, bool]:
|
||||||
"""Grid-walk layer_id в bbox квартала.
|
"""Grid-walk layer_id в bbox квартала.
|
||||||
|
|
||||||
Генерирует grid_size×grid_size ячеек, для каждой делает wms_feature_info.
|
Генерирует grid_size×grid_size ячеек, для каждой делает wms_feature_info.
|
||||||
|
|
@ -306,7 +324,14 @@ async def _grid_walk_category(
|
||||||
помечал worker как мёртвого. Вызываем каждые heartbeat_every cells.
|
помечал worker как мёртвого. Вызываем каждые heartbeat_every cells.
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
(upserted_count, requests_count)
|
(upserted_count, requests_count, layer_failed).
|
||||||
|
|
||||||
|
layer_failed=True — issue #252 — когда ВСЕ выполненные запросы слоя
|
||||||
|
упали на server-side 500/ServiceException и ни один не прошёл. Это
|
||||||
|
отличает «слой реально пуст в квартале» (0 discovered, 0 ошибок) от
|
||||||
|
«слой временно недоступен на стороне NSPD» (0 discovered потому что
|
||||||
|
каждый cell отдал 500). Вызывающий фиксирует флаг в harvest_meta —
|
||||||
|
квартал НЕ валится, один сбойный слой просто пропускается.
|
||||||
"""
|
"""
|
||||||
bbox = quarter_bbox_3857(db, quarter)
|
bbox = quarter_bbox_3857(db, quarter)
|
||||||
if bbox is None:
|
if bbox is None:
|
||||||
|
|
@ -315,13 +340,15 @@ async def _grid_walk_category(
|
||||||
quarter,
|
quarter,
|
||||||
layer_id,
|
layer_id,
|
||||||
)
|
)
|
||||||
return 0, 0
|
return 0, 0, False
|
||||||
|
|
||||||
grid_points = generate_grid_click_points(bbox, grid_size=grid_size, tile_size=tile_size)
|
grid_points = generate_grid_click_points(bbox, grid_size=grid_size, tile_size=tile_size)
|
||||||
|
|
||||||
discovered_cads: set[str] = set()
|
discovered_cads: set[str] = set()
|
||||||
upserted = 0
|
upserted = 0
|
||||||
requests = 0
|
requests = 0
|
||||||
|
server_errors = 0 # cells, упавшие на 5xx/ServiceException (issue #252)
|
||||||
|
ok_cells = 0 # cells, вернувшие ответ (пусть и пустой) — слой жив
|
||||||
|
|
||||||
for idx, (cell_bbox, click_xy) in enumerate(grid_points):
|
for idx, (cell_bbox, click_xy) in enumerate(grid_points):
|
||||||
# Fix I: heartbeat каждые N cells — обновляет cadastre_jobs.heartbeat_at
|
# Fix I: heartbeat каждые N cells — обновляет cadastre_jobs.heartbeat_at
|
||||||
|
|
@ -345,14 +372,31 @@ async def _grid_walk_category(
|
||||||
height=tile_size,
|
height=tile_size,
|
||||||
)
|
)
|
||||||
requests += 1
|
requests += 1
|
||||||
except Exception as e:
|
ok_cells += 1
|
||||||
# Fix J: NSPD WMS возвращает HTTP 500 на десятки cells per quarter —
|
except NspdBulkServerError as e:
|
||||||
# это server-side noise, не наш bug. logger.debug чтобы не засорять
|
# Issue #252: 5xx / WMS ServiceException — transient server-side noise
|
||||||
# prod логи (раньше было warning → 200+ warnings per quarter).
|
# NSPD (десятки cells per quarter). logger.debug чтобы не засорять prod
|
||||||
|
# логи. Считаем отдельно: если ВСЕ cells слоя упали так и ни один не
|
||||||
|
# прошёл — поднимем layer_failed (см. return) → harvest_meta-флаг.
|
||||||
logger.debug(
|
logger.debug(
|
||||||
"_grid_walk_category: wms_feature_info error layer=%d quarter=%s: %s",
|
"_grid_walk_category: server error layer=%d quarter=%s cell=%d: %s",
|
||||||
layer_id,
|
layer_id,
|
||||||
quarter,
|
quarter,
|
||||||
|
idx,
|
||||||
|
e,
|
||||||
|
)
|
||||||
|
requests += 1
|
||||||
|
server_errors += 1
|
||||||
|
continue
|
||||||
|
except Exception as e:
|
||||||
|
# Прочие (сетевые / parse) ошибки одного cell — тоже не валим квартал,
|
||||||
|
# но это НЕ server-side 500 → не учитываем в server_errors (иначе сеть
|
||||||
|
# ложно triggers layer_failed). Слой жив, просто этот cell не дошёл.
|
||||||
|
logger.debug(
|
||||||
|
"_grid_walk_category: wms_feature_info error layer=%d quarter=%s cell=%d: %s",
|
||||||
|
layer_id,
|
||||||
|
quarter,
|
||||||
|
idx,
|
||||||
e,
|
e,
|
||||||
)
|
)
|
||||||
requests += 1
|
requests += 1
|
||||||
|
|
@ -376,7 +420,190 @@ async def _grid_walk_category(
|
||||||
e,
|
e,
|
||||||
)
|
)
|
||||||
|
|
||||||
return upserted, requests
|
# Issue #252: слой считаем «сбойным» только если БЫЛИ server-side ошибки И
|
||||||
|
# ни один cell не прошёл успешно. ok_cells>0 → слой жив (0 discovered значит
|
||||||
|
# реально пусто). server_errors==0 → штатный пустой/полный обход.
|
||||||
|
layer_failed = server_errors > 0 and ok_cells == 0
|
||||||
|
if layer_failed:
|
||||||
|
logger.warning(
|
||||||
|
"_grid_walk_category: layer=%d quarter=%s ПОЛНОСТЬЮ сбойный — "
|
||||||
|
"%d/%d cells вернули 5xx/ServiceException, 0 успешных. Слой пропущен.",
|
||||||
|
layer_id,
|
||||||
|
quarter,
|
||||||
|
server_errors,
|
||||||
|
len(grid_points),
|
||||||
|
)
|
||||||
|
|
||||||
|
return upserted, requests, layer_failed
|
||||||
|
|
||||||
|
|
||||||
|
# ── Geom backfill для участков с geom IS NULL (issue #200) ────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass
|
||||||
|
class GeomBackfillResult:
|
||||||
|
"""Итог докачки geom для участков с geom IS NULL (issue #200)."""
|
||||||
|
|
||||||
|
quarters_scanned: int = 0
|
||||||
|
parcels_targeted: int = 0 # участки без geom на входе
|
||||||
|
parcels_healed: int = 0 # получили полигон после grid-walk
|
||||||
|
parcels_marked_unavailable: int = 0 # geom так и не нашёлся → флаг
|
||||||
|
grid_walk_requests: int = 0
|
||||||
|
|
||||||
|
|
||||||
|
async def backfill_parcel_geom(
|
||||||
|
db: Session,
|
||||||
|
client: NSPDBulkClient,
|
||||||
|
*,
|
||||||
|
limit: int = 500,
|
||||||
|
update_heartbeat: Callable[[dict[str, Any]], None] | None = None,
|
||||||
|
) -> GeomBackfillResult:
|
||||||
|
"""Докачать geom для участков cad_parcels с geom IS NULL (issue #200).
|
||||||
|
|
||||||
|
Алгоритм:
|
||||||
|
1. Выбрать до `limit` участков `geom IS NULL AND NOT geom_unavailable`.
|
||||||
|
2. Сгруппировать по кварталу (3-сегментный cad). NSPD search-эндпоинт
|
||||||
|
(Phase 1) отдаёт эти участки без полигона; реальную геометрию
|
||||||
|
возвращает grid-walk слоя 36368 (WMS GetFeatureInfo по bbox квартала).
|
||||||
|
3. Для каждого квартала вызвать _grid_walk_category(36368) — он upsert'ит
|
||||||
|
полигоны через ON CONFLICT ... geom = COALESCE(EXCLUDED.geom, ...).
|
||||||
|
4. Участки, которые ПОСЛЕ grid-walk всё ещё geom IS NULL, пометить
|
||||||
|
geom_unavailable=TRUE — geom систематически недоступен (кросс-региональные
|
||||||
|
артефакты поиска без полигона в слое), чтобы не ретраить вечно.
|
||||||
|
|
||||||
|
Идемпотентен: повторный вызов берёт следующую порцию (healed уходят из
|
||||||
|
выборки — у них geom уже не NULL; unavailable исключены флагом).
|
||||||
|
|
||||||
|
Args:
|
||||||
|
limit: максимум участков за один прогон (батч против WAF-burst).
|
||||||
|
update_heartbeat: callback для heartbeat (как в harvest_quarter) — задача
|
||||||
|
может идти минуты при многих кварталах.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
GeomBackfillResult со счётчиками.
|
||||||
|
"""
|
||||||
|
result = GeomBackfillResult()
|
||||||
|
|
||||||
|
rows = (
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"SELECT cad_num FROM cad_parcels "
|
||||||
|
"WHERE geom IS NULL AND geom_unavailable = FALSE "
|
||||||
|
"ORDER BY cad_num LIMIT :lim"
|
||||||
|
),
|
||||||
|
{"lim": limit},
|
||||||
|
)
|
||||||
|
.scalars()
|
||||||
|
.all()
|
||||||
|
)
|
||||||
|
if not rows:
|
||||||
|
logger.info("backfill_parcel_geom: нет участков с geom IS NULL — нечего докачивать")
|
||||||
|
return result
|
||||||
|
|
||||||
|
result.parcels_targeted = len(rows)
|
||||||
|
|
||||||
|
# Группируем по кварталу (первые 3 colon-сегмента). Участки без валидного
|
||||||
|
# 3-сегментного префикса (мусорные cad) — сразу в unavailable, grid-walk
|
||||||
|
# для них невозможен (нет квартала → нет bbox).
|
||||||
|
by_quarter: dict[str, list[str]] = {}
|
||||||
|
orphan_cads: list[str] = []
|
||||||
|
for cad in rows:
|
||||||
|
parts = cad.split(":")
|
||||||
|
if len(parts) >= 3 and all(parts[:3]):
|
||||||
|
by_quarter.setdefault(":".join(parts[:3]), []).append(cad)
|
||||||
|
else:
|
||||||
|
orphan_cads.append(cad)
|
||||||
|
|
||||||
|
if orphan_cads:
|
||||||
|
result.parcels_marked_unavailable += _mark_geom_unavailable(db, orphan_cads)
|
||||||
|
db.commit()
|
||||||
|
|
||||||
|
for quarter, cads in by_quarter.items():
|
||||||
|
result.quarters_scanned += 1
|
||||||
|
if update_heartbeat is not None:
|
||||||
|
update_heartbeat(
|
||||||
|
{
|
||||||
|
"phase": "geom_backfill",
|
||||||
|
"quarter": quarter,
|
||||||
|
"quarters_scanned": result.quarters_scanned,
|
||||||
|
"quarters_total": len(by_quarter),
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
try:
|
||||||
|
_discovered, n_requests, _failed = await _grid_walk_category(
|
||||||
|
db=db,
|
||||||
|
client=client,
|
||||||
|
quarter=quarter,
|
||||||
|
layer_id=CAT_PARCEL,
|
||||||
|
update_progress=update_heartbeat,
|
||||||
|
)
|
||||||
|
result.grid_walk_requests += n_requests
|
||||||
|
db.commit()
|
||||||
|
except Exception as e:
|
||||||
|
# Один сбойный квартал не валит весь backfill — лог + продолжаем.
|
||||||
|
# (WAF 403 пробросится из client и прервёт прогон — это ожидаемо,
|
||||||
|
# caller-task ловит и не ретраит, как в bulk_harvest.)
|
||||||
|
logger.warning(
|
||||||
|
"backfill_parcel_geom: grid-walk failed quarter=%s: %s", quarter, e
|
||||||
|
)
|
||||||
|
db.rollback()
|
||||||
|
continue
|
||||||
|
|
||||||
|
# Какие из целевых участков квартала всё ещё без geom?
|
||||||
|
still_null = (
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"SELECT cad_num FROM cad_parcels "
|
||||||
|
"WHERE cad_num = ANY(CAST(:cads AS text[])) AND geom IS NULL"
|
||||||
|
),
|
||||||
|
{"cads": cads},
|
||||||
|
)
|
||||||
|
.scalars()
|
||||||
|
.all()
|
||||||
|
)
|
||||||
|
healed_here = len(cads) - len(still_null)
|
||||||
|
result.parcels_healed += healed_here
|
||||||
|
|
||||||
|
if still_null:
|
||||||
|
result.parcels_marked_unavailable += _mark_geom_unavailable(db, list(still_null))
|
||||||
|
db.commit()
|
||||||
|
|
||||||
|
logger.info(
|
||||||
|
"backfill_parcel_geom: quarter=%s targeted=%d healed=%d unavailable=%d",
|
||||||
|
quarter,
|
||||||
|
len(cads),
|
||||||
|
healed_here,
|
||||||
|
len(still_null),
|
||||||
|
)
|
||||||
|
|
||||||
|
logger.info(
|
||||||
|
"backfill_parcel_geom done: quarters=%d targeted=%d healed=%d unavailable=%d requests=%d",
|
||||||
|
result.quarters_scanned,
|
||||||
|
result.parcels_targeted,
|
||||||
|
result.parcels_healed,
|
||||||
|
result.parcels_marked_unavailable,
|
||||||
|
result.grid_walk_requests,
|
||||||
|
)
|
||||||
|
return result
|
||||||
|
|
||||||
|
|
||||||
|
def _mark_geom_unavailable(db: Session, cad_nums: list[str]) -> int:
|
||||||
|
"""Пометить участки geom_unavailable=TRUE (issue #200). Возвращает rowcount.
|
||||||
|
|
||||||
|
Параметризованный UPDATE по списку cad_num (ANY(text[])) — не f-string SQL.
|
||||||
|
"""
|
||||||
|
if not cad_nums:
|
||||||
|
return 0
|
||||||
|
res = db.execute(
|
||||||
|
text(
|
||||||
|
"UPDATE cad_parcels SET geom_unavailable = TRUE, updated_at = NOW() "
|
||||||
|
"WHERE cad_num = ANY(CAST(:cads AS text[])) "
|
||||||
|
"AND geom IS NULL AND geom_unavailable = FALSE"
|
||||||
|
),
|
||||||
|
{"cads": cad_nums},
|
||||||
|
)
|
||||||
|
return res.rowcount or 0
|
||||||
|
|
||||||
|
|
||||||
# ── Routing dispatcher ───────────────────────────────────────────────────────
|
# ── Routing dispatcher ───────────────────────────────────────────────────────
|
||||||
|
|
|
||||||
|
|
@ -277,17 +277,23 @@ def build_beat_schedule() -> dict:
|
||||||
"options": {"queue": "celery"},
|
"options": {"queue": "celery"},
|
||||||
}
|
}
|
||||||
|
|
||||||
# ПЗЗ территориальных зон ЕКБ из PKK6 ArcGIS — ежемесячно 1-го числа в 03:00 МСК.
|
# ПЗЗ территориальных зон ЕКБ из PKK6 ArcGIS — DISABLED 2026-06-13 (issue #259).
|
||||||
# Celery conf.timezone=Europe/Moscow → crontab трактуется в МСК (#1233).
|
#
|
||||||
# PKK6 нестабилен под нагрузкой, но данные ПЗЗ меняются редко — раз в месяц достаточно.
|
# ПОЧЕМУ ОТКЛЮЧЕНО: источник pkk.rosreestr.ru/.../PKK6/ZONES deprecated (Росреестр
|
||||||
# Task: tasks/pzz_sync.py → sync_pzz_zones_ekb.
|
# вывел PKK6-эндпоинт), таблица pzz_zones_ekb пуста (0 rows на prod) — задача ни
|
||||||
# Admin trigger: POST /api/v1/admin/scrape/pzz-sync.
|
# разу не наполнила её успешно. ПЗЗ-данные пришли в систему ИНЫМ путём:
|
||||||
# Ref: issue #233 (pzz_zones_ekb = 0 rows, задача никогда не запускалась автоматически).
|
# zone_regulation_cache (#1059, beat zone-regulation-refresh-monthly) + NSPD
|
||||||
schedule["pzz-sync-monthly"] = {
|
# territorial_zones dumps (Phase 2.5 bulk_harvest, layer 875838). Активный beat
|
||||||
"task": "tasks.pzz_sync.sync_pzz_zones_ekb",
|
# каждый месяц дёргал deprecated PKK6 → broad except в pzz_sync.py логировал
|
||||||
"schedule": _parse_cron("0 3 1 * *"),
|
# ошибку как error → рекуррентил GlitchTip BACKEND-1B.
|
||||||
"options": {"queue": "celery"},
|
#
|
||||||
}
|
# Ручной перезапуск (если появится рабочий источник) остаётся через admin-эндпоинт
|
||||||
|
# POST /api/v1/admin/scrape/pzz-sync. Ref: #233 (исходно «0 rows»), #259 (disable).
|
||||||
|
# schedule["pzz-sync-monthly"] = {
|
||||||
|
# "task": "tasks.pzz_sync.sync_pzz_zones_ekb",
|
||||||
|
# "schedule": _parse_cron("0 3 1 * *"),
|
||||||
|
# "options": {"queue": "celery"},
|
||||||
|
# }
|
||||||
|
|
||||||
# Catalog-object scrape — наполняет ~25 NULL колонок domrf_kn_objects из SSR-страниц.
|
# Catalog-object scrape — наполняет ~25 NULL колонок domrf_kn_objects из SSR-страниц.
|
||||||
# kn-API не отдаёт wall_type, energy_eff, ceiling_height_m, parking_* и т.д.
|
# kn-API не отдаёт wall_type, energy_eff, ceiling_height_m, parking_* и т.д.
|
||||||
|
|
|
||||||
|
|
@ -34,6 +34,12 @@ def _log_breadcrumb(level: str, message: str) -> None:
|
||||||
def sync_pzz_zones_ekb() -> dict:
|
def sync_pzz_zones_ekb() -> dict:
|
||||||
"""Импорт ПЗЗ территориальных зон ЕКБ из Росреестр PKK6.
|
"""Импорт ПЗЗ территориальных зон ЕКБ из Росреестр PKK6.
|
||||||
|
|
||||||
|
DEPRECATED (issue #259): источник PKK6 выведен Росреестром, beat-расписание
|
||||||
|
pzz-sync-monthly отключено (beat_schedule.py), таблица pzz_zones_ekb пуста.
|
||||||
|
ПЗЗ-данные теперь поступают через zone_regulation_cache + NSPD territorial_zones
|
||||||
|
dumps. Задача оставлена только для ручного запуска через admin (на случай
|
||||||
|
появления рабочего источника).
|
||||||
|
|
||||||
Запускается вручную через POST /api/v1/admin/scrape/pzz-sync.
|
Запускается вручную через POST /api/v1/admin/scrape/pzz-sync.
|
||||||
Persistent breadcrumbs в nspd_geo_log (stage='pzz_sync') для диагностики.
|
Persistent breadcrumbs в nspd_geo_log (stage='pzz_sync') для диагностики.
|
||||||
"""
|
"""
|
||||||
|
|
@ -46,6 +52,10 @@ def sync_pzz_zones_ekb() -> dict:
|
||||||
_log_breadcrumb("info", f"done: {result}")
|
_log_breadcrumb("info", f"done: {result}")
|
||||||
return result
|
return result
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception("pzz_sync failed: %s", e)
|
# Issue #259: PKK6 deprecated → этот путь почти всегда падает. logger.warning
|
||||||
_log_breadcrumb("error", f"{type(e).__name__}: {e}")
|
# (не .exception) — Sentry/GlitchTip НЕ захватывает warning как error-event,
|
||||||
|
# поэтому ручной запуск по дохлому источнику не рекуррентит BACKEND-1B.
|
||||||
|
# re-raise сохранён: admin-эндпоинт обязан увидеть фактический провал импорта.
|
||||||
|
logger.warning("pzz_sync failed (PKK6 deprecated, см. #259): %s: %s", type(e).__name__, e)
|
||||||
|
_log_breadcrumb("warning", f"{type(e).__name__}: {e}")
|
||||||
raise
|
raise
|
||||||
|
|
|
||||||
|
|
@ -24,7 +24,7 @@ from typing import Any
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
|
|
||||||
from app.core.db import SessionLocal
|
from app.core.db import SessionLocal
|
||||||
from app.scrapers.nspd_bulk_client import NspdBulkWafError
|
from app.scrapers.nspd_bulk_client import NspdBulkRateLimitError, NspdBulkWafError
|
||||||
from app.workers.celery_app import celery_app
|
from app.workers.celery_app import celery_app
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
@ -132,6 +132,12 @@ def _is_job_cancelled(db: Any, job_id: int) -> bool:
|
||||||
name="tasks.cadastre.harvest_quarter",
|
name="tasks.cadastre.harvest_quarter",
|
||||||
acks_late=True,
|
acks_late=True,
|
||||||
max_retries=3,
|
max_retries=3,
|
||||||
|
# autoretry_for=(Exception,) ловит ВСЁ кроме явных исключений ниже. Issue #252:
|
||||||
|
# квартал-уровневый NspdBulkServerError (HTTP 500 / WMS ServiceException на
|
||||||
|
# Phase 1 search_by_quarter) — подкласс Exception и НЕ в dont_autoretry_for,
|
||||||
|
# значит ретраится с backoff (GlitchTip BACKEND-8/10/11/.../37). Per-cell 500
|
||||||
|
# внутри grid-walk квартал НЕ валят — слой пропускается с harvest_meta-флагом
|
||||||
|
# (см. bulk_harvest._grid_walk_category).
|
||||||
autoretry_for=(Exception,),
|
autoretry_for=(Exception,),
|
||||||
retry_backoff=True,
|
retry_backoff=True,
|
||||||
retry_backoff_max=60,
|
retry_backoff_max=60,
|
||||||
|
|
@ -211,6 +217,9 @@ def bulk_harvest_quarter_task(self: Any, quarter: str, job_id: int) -> dict[str,
|
||||||
"snapshot_requests": result.snapshot_requests,
|
"snapshot_requests": result.snapshot_requests,
|
||||||
"grid_walk_requests": result.grid_walk_requests,
|
"grid_walk_requests": result.grid_walk_requests,
|
||||||
"total_requests": result.total_requests,
|
"total_requests": result.total_requests,
|
||||||
|
# Issue #252: слои, чей grid-walk целиком сбойнул на 500 — для
|
||||||
|
# observability в результате task (квартал при этом успешен).
|
||||||
|
"failed_layers": result.failed_layers,
|
||||||
}
|
}
|
||||||
|
|
||||||
result_dict = asyncio.run(_run())
|
result_dict = asyncio.run(_run())
|
||||||
|
|
@ -472,3 +481,74 @@ def cleanup_cadastre_zombies() -> dict[str, Any]:
|
||||||
finally:
|
finally:
|
||||||
db.close()
|
db.close()
|
||||||
return {"stale_paused": stale_ids}
|
return {"stale_paused": stale_ids}
|
||||||
|
|
||||||
|
|
||||||
|
# ── Geom backfill (issue #200) ───────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
@celery_app.task(
|
||||||
|
bind=True,
|
||||||
|
name="tasks.cadastre.backfill_parcel_geom",
|
||||||
|
acks_late=True,
|
||||||
|
max_retries=2,
|
||||||
|
# WAF/429 — не наш баг, не ретраим (как bulk_harvest_quarter_task). Прочие
|
||||||
|
# ошибки докачки тоже не ретраим автоматически: backfill идемпотентен и
|
||||||
|
# запускается вручную/по расписанию, повторный прогон возьмёт остаток.
|
||||||
|
autoretry_for=(),
|
||||||
|
soft_time_limit=900,
|
||||||
|
time_limit=1200,
|
||||||
|
)
|
||||||
|
def backfill_parcel_geom_task(self: Any, limit: int = 500) -> dict[str, Any]:
|
||||||
|
"""Докачать geom для участков cad_parcels с geom IS NULL (issue #200).
|
||||||
|
|
||||||
|
Sync→async bridge через asyncio.run() (как bulk_harvest_quarter_task).
|
||||||
|
Один прогон обрабатывает до `limit` участков; idempotent — повторный вызов
|
||||||
|
берёт следующую порцию (вылеченные уходят из выборки, систематически
|
||||||
|
недоступные исключаются флагом geom_unavailable).
|
||||||
|
|
||||||
|
Запуск вручную (CLI/admin) или через beat, если нужно периодически дочищать
|
||||||
|
хвост NULL-geom после bulk-harvest. Возвращает счётчики GeomBackfillResult.
|
||||||
|
"""
|
||||||
|
from app.services.cadastre.bulk_harvest import backfill_parcel_geom
|
||||||
|
|
||||||
|
db = SessionLocal()
|
||||||
|
try:
|
||||||
|
|
||||||
|
async def _run() -> dict[str, Any]:
|
||||||
|
from app.scrapers.nspd_bulk_client import NSPDBulkClient
|
||||||
|
|
||||||
|
async with NSPDBulkClient() as client:
|
||||||
|
res = await backfill_parcel_geom(db=db, client=client, limit=limit)
|
||||||
|
return {
|
||||||
|
"quarters_scanned": res.quarters_scanned,
|
||||||
|
"parcels_targeted": res.parcels_targeted,
|
||||||
|
"parcels_healed": res.parcels_healed,
|
||||||
|
"parcels_marked_unavailable": res.parcels_marked_unavailable,
|
||||||
|
"grid_walk_requests": res.grid_walk_requests,
|
||||||
|
}
|
||||||
|
|
||||||
|
result_dict = asyncio.run(_run())
|
||||||
|
logger.info("backfill_parcel_geom_task done: %s", result_dict)
|
||||||
|
return result_dict
|
||||||
|
|
||||||
|
except (NspdBulkWafError, NspdBulkRateLimitError) as e:
|
||||||
|
# WAF 403 / исчерпанный rate-limit — прерываем прогон без retry. Уже
|
||||||
|
# вылеченные участки закоммичены поквартально внутри backfill_parcel_geom,
|
||||||
|
# следующий ручной/плановый запуск добьёт остаток.
|
||||||
|
logger.warning("backfill_parcel_geom_task aborted (WAF/rate-limit): %s", e)
|
||||||
|
try:
|
||||||
|
db.rollback()
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
return {"aborted": "waf_or_rate_limit", "error": str(e)[:300]}
|
||||||
|
|
||||||
|
except Exception as e:
|
||||||
|
logger.exception("backfill_parcel_geom_task failed: %s", e)
|
||||||
|
try:
|
||||||
|
db.rollback()
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
raise
|
||||||
|
|
||||||
|
finally:
|
||||||
|
db.close()
|
||||||
|
|
|
||||||
|
|
@ -410,6 +410,75 @@ async def test_http_429_backoff_exhausted() -> None:
|
||||||
assert call_count == 4
|
assert call_count == 4
|
||||||
|
|
||||||
|
|
||||||
|
# ── Issue #252: 5xx / WMS ServiceException → NspdBulkServerError ──────────────
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_http_500_raises_server_error() -> None:
|
||||||
|
"""500 → NspdBulkServerError (подкласс NspdBulkError) — для retry + skip."""
|
||||||
|
import httpx
|
||||||
|
|
||||||
|
from app.scrapers.nspd_bulk_client import NspdBulkError, NspdBulkServerError
|
||||||
|
|
||||||
|
async def fake_get(*args: Any, **kwargs: Any) -> httpx.Response:
|
||||||
|
return httpx.Response(500, text="Internal Server Error", request=MagicMock())
|
||||||
|
|
||||||
|
async with NSPDBulkClient() as client:
|
||||||
|
assert client._client is not None
|
||||||
|
with patch.object(client._client, "get", new=AsyncMock(side_effect=fake_get)):
|
||||||
|
with pytest.raises(NspdBulkServerError) as exc_info:
|
||||||
|
await client._get_json("https://nspd.gov.ru/test")
|
||||||
|
|
||||||
|
# Подкласс NspdBulkError → autoretry_for=(Exception,) его поймает, а
|
||||||
|
# except NspdBulkError-ветки (404/400) тоже сработают если понадобится.
|
||||||
|
assert isinstance(exc_info.value, NspdBulkError)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_wms_service_exception_in_200_body_raises_server_error() -> None:
|
||||||
|
"""HTTP 200 с XML <ServiceException> в теле → NspdBulkServerError (не JSONDecodeError)."""
|
||||||
|
import httpx
|
||||||
|
|
||||||
|
from app.scrapers.nspd_bulk_client import NspdBulkServerError
|
||||||
|
|
||||||
|
xml_body = (
|
||||||
|
'<?xml version="1.0"?><ServiceExceptionReport>'
|
||||||
|
"<ServiceException>Rendering error</ServiceException></ServiceExceptionReport>"
|
||||||
|
)
|
||||||
|
|
||||||
|
async def fake_get(*args: Any, **kwargs: Any) -> httpx.Response:
|
||||||
|
return httpx.Response(
|
||||||
|
200, text=xml_body, headers={"content-type": "application/xml"}, request=MagicMock()
|
||||||
|
)
|
||||||
|
|
||||||
|
async with NSPDBulkClient() as client:
|
||||||
|
assert client._client is not None
|
||||||
|
with patch.object(client._client, "get", new=AsyncMock(side_effect=fake_get)):
|
||||||
|
with pytest.raises(NspdBulkServerError):
|
||||||
|
await client._get_json("https://nspd.gov.ru/wms")
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_wms_feature_info_propagates_server_error(
|
||||||
|
sample_wms_response: dict[str, Any],
|
||||||
|
) -> None:
|
||||||
|
"""wms_feature_info пробрасывает NspdBulkServerError (caller считает fail-rate)."""
|
||||||
|
from app.scrapers.nspd_bulk_client import NspdBulkServerError
|
||||||
|
|
||||||
|
async with NSPDBulkClient() as client:
|
||||||
|
with patch.object(
|
||||||
|
client,
|
||||||
|
"_get_json",
|
||||||
|
new=AsyncMock(side_effect=NspdBulkServerError("HTTP 500 server error")),
|
||||||
|
):
|
||||||
|
with pytest.raises(NspdBulkServerError):
|
||||||
|
await client.wms_feature_info(
|
||||||
|
layer_id=36368,
|
||||||
|
bbox=(6700000.0, 7700000.0, 6701000.0, 7701000.0),
|
||||||
|
click_xy=(256, 256),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
# ── E2E tests (network, skipped by default) ───────────────────────────────────
|
# ── E2E tests (network, skipped by default) ───────────────────────────────────
|
||||||
|
|
||||||
pytestmark_slow = pytest.mark.slow
|
pytestmark_slow = pytest.mark.slow
|
||||||
|
|
|
||||||
|
|
@ -673,7 +673,7 @@ async def test_grid_walk_dedup_same_cad_num() -> None:
|
||||||
}
|
}
|
||||||
|
|
||||||
# Используем grid_size=2 для скорости теста (4 ячейки)
|
# Используем grid_size=2 для скорости теста (4 ячейки)
|
||||||
upserted, requests = await _grid_walk_category(
|
upserted, requests, layer_failed = await _grid_walk_category(
|
||||||
db=db,
|
db=db,
|
||||||
client=client,
|
client=client,
|
||||||
quarter="66:41:0303161",
|
quarter="66:41:0303161",
|
||||||
|
|
@ -685,6 +685,7 @@ async def test_grid_walk_dedup_same_cad_num() -> None:
|
||||||
assert mock_upsert.call_count == 1
|
assert mock_upsert.call_count == 1
|
||||||
assert upserted == 1
|
assert upserted == 1
|
||||||
assert requests == 4 # 2×2 grid
|
assert requests == 4 # 2×2 grid
|
||||||
|
assert layer_failed is False # все cells успешны
|
||||||
|
|
||||||
|
|
||||||
# ── ЗОУИТ dedup по (reg_numb_border, category_id) ───────────────────────────
|
# ── ЗОУИТ dedup по (reg_numb_border, category_id) ───────────────────────────
|
||||||
|
|
@ -1306,7 +1307,7 @@ async def test_grid_walk_emits_heartbeat_callbacks() -> None:
|
||||||
|
|
||||||
progress_states: list[dict[str, Any]] = []
|
progress_states: list[dict[str, Any]] = []
|
||||||
|
|
||||||
_upserted, requests = await _grid_walk_category(
|
_upserted, requests, _failed = await _grid_walk_category(
|
||||||
db=db,
|
db=db,
|
||||||
client=client,
|
client=client,
|
||||||
quarter="66:41:0303161",
|
quarter="66:41:0303161",
|
||||||
|
|
@ -1341,10 +1342,297 @@ async def test_grid_walk_no_heartbeat_when_callback_none() -> None:
|
||||||
client.wms_feature_info = AsyncMock(return_value=[])
|
client.wms_feature_info = AsyncMock(return_value=[])
|
||||||
|
|
||||||
# Не передаём update_progress — должно не падать
|
# Не передаём update_progress — должно не падать
|
||||||
_upserted, requests = await _grid_walk_category(
|
_upserted, requests, _failed = await _grid_walk_category(
|
||||||
db=db,
|
db=db,
|
||||||
client=client,
|
client=client,
|
||||||
quarter="66:41:0303161",
|
quarter="66:41:0303161",
|
||||||
layer_id=36368,
|
layer_id=36368,
|
||||||
)
|
)
|
||||||
assert requests == 225
|
assert requests == 225
|
||||||
|
|
||||||
|
|
||||||
|
# ── Issue #252: per-layer skip при 5xx/ServiceException ──────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def _mock_db_grid_bbox() -> MagicMock:
|
||||||
|
"""Mock db с валидным bbox для quarter_bbox_3857 + begin_nested CM."""
|
||||||
|
db = MagicMock()
|
||||||
|
db.execute.return_value.mappings.return_value.first.return_value = {
|
||||||
|
"xmin": 6735845.0,
|
||||||
|
"ymin": 8329000.0,
|
||||||
|
"xmax": 6736595.0,
|
||||||
|
"ymax": 8329750.0,
|
||||||
|
}
|
||||||
|
db.begin_nested = MagicMock(
|
||||||
|
return_value=MagicMock(__enter__=MagicMock(return_value=None), __exit__=MagicMock())
|
||||||
|
)
|
||||||
|
db.commit = MagicMock()
|
||||||
|
return db
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_grid_walk_marks_layer_failed_when_all_cells_500() -> None:
|
||||||
|
"""Issue #252: ВСЕ cells слоя упали на NspdBulkServerError → layer_failed=True."""
|
||||||
|
from app.scrapers.nspd_bulk_client import NspdBulkServerError
|
||||||
|
from app.services.cadastre.bulk_harvest import _grid_walk_category
|
||||||
|
|
||||||
|
db = _mock_db_grid_bbox()
|
||||||
|
client = AsyncMock()
|
||||||
|
client.wms_feature_info = AsyncMock(side_effect=NspdBulkServerError("HTTP 500"))
|
||||||
|
|
||||||
|
upserted, requests, layer_failed = await _grid_walk_category(
|
||||||
|
db=db,
|
||||||
|
client=client,
|
||||||
|
quarter="66:41:0303161",
|
||||||
|
layer_id=36368,
|
||||||
|
grid_size=3, # 9 cells, все упадут
|
||||||
|
)
|
||||||
|
|
||||||
|
assert upserted == 0
|
||||||
|
assert requests == 9
|
||||||
|
assert layer_failed is True
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_grid_walk_layer_not_failed_when_some_cells_ok() -> None:
|
||||||
|
"""Issue #252: если хоть один cell прошёл — layer_failed=False (слой жив, просто пуст)."""
|
||||||
|
from app.scrapers.nspd_bulk_client import NspdBulkServerError
|
||||||
|
from app.services.cadastre.bulk_harvest import _grid_walk_category
|
||||||
|
|
||||||
|
db = _mock_db_grid_bbox()
|
||||||
|
client = AsyncMock()
|
||||||
|
# Первый cell ок (пустой ответ), остальные 500
|
||||||
|
client.wms_feature_info = AsyncMock(
|
||||||
|
side_effect=[[], NspdBulkServerError("500"), NspdBulkServerError("500")]
|
||||||
|
)
|
||||||
|
|
||||||
|
_upserted, requests, layer_failed = await _grid_walk_category(
|
||||||
|
db=db,
|
||||||
|
client=client,
|
||||||
|
quarter="66:41:0303161",
|
||||||
|
layer_id=36368,
|
||||||
|
grid_size=1, # уйдёт по side_effect — но grid_size=1 → 1 cell
|
||||||
|
)
|
||||||
|
|
||||||
|
# grid_size=1 → 1 cell (ок) → layer жив
|
||||||
|
assert layer_failed is False
|
||||||
|
assert requests == 1
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_grid_walk_network_error_does_not_trigger_layer_failed() -> None:
|
||||||
|
"""Issue #252: сетевая (не server-side) ошибка cell НЕ учитывается в fail-rate."""
|
||||||
|
from app.scrapers.nspd_bulk_client import NspdBulkError
|
||||||
|
from app.services.cadastre.bulk_harvest import _grid_walk_category
|
||||||
|
|
||||||
|
db = _mock_db_grid_bbox()
|
||||||
|
client = AsyncMock()
|
||||||
|
# NspdBulkError (не Server) — например network. layer_failed требует server_errors>0.
|
||||||
|
client.wms_feature_info = AsyncMock(side_effect=NspdBulkError("Network error"))
|
||||||
|
|
||||||
|
_upserted, requests, layer_failed = await _grid_walk_category(
|
||||||
|
db=db,
|
||||||
|
client=client,
|
||||||
|
quarter="66:41:0303161",
|
||||||
|
layer_id=36368,
|
||||||
|
grid_size=3,
|
||||||
|
)
|
||||||
|
|
||||||
|
# server_errors==0 (это не 5xx) → layer_failed=False, квартал не получит skip-флаг
|
||||||
|
assert layer_failed is False
|
||||||
|
assert requests == 9
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_harvest_quarter_records_failed_layer_in_phase_state() -> None:
|
||||||
|
"""Issue #252: сбойный слой grid-walk → harvest_meta.layer_X_failed в done-state."""
|
||||||
|
from app.services.cadastre.bulk_harvest import CAT_PARCEL, harvest_quarter
|
||||||
|
|
||||||
|
parcel = _make_parcel_feature()
|
||||||
|
snapshot = QuarterSnapshot(
|
||||||
|
quarter="66:41:0303161",
|
||||||
|
fetched_at="2026-05-15T10:00:00+00:00",
|
||||||
|
features=[parcel],
|
||||||
|
meta_counts={CAT_PARCEL: 170}, # overflow → grid-walk слоя 36368
|
||||||
|
)
|
||||||
|
|
||||||
|
db = _mock_db_grid_bbox()
|
||||||
|
client = AsyncMock()
|
||||||
|
client.search_by_quarter = AsyncMock(return_value=snapshot)
|
||||||
|
client.get_territorial_zones_in_bbox = AsyncMock(return_value=[])
|
||||||
|
|
||||||
|
progress_states: list[dict[str, Any]] = []
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch("app.services.cadastre.bulk_harvest.upsert_features") as mock_upsert,
|
||||||
|
patch("app.services.cadastre.bulk_harvest._grid_walk_category") as mock_walk,
|
||||||
|
patch("app.services.cadastre.bulk_harvest.quarter_bbox_3857") as mock_bbox,
|
||||||
|
):
|
||||||
|
mock_bbox.return_value = (6735845.0, 8329000.0, 6736595.0, 8329750.0)
|
||||||
|
mock_upsert.return_value = {
|
||||||
|
"parcels": 1,
|
||||||
|
"buildings": 0,
|
||||||
|
"constructions": 0,
|
||||||
|
"oncs": 0,
|
||||||
|
"enks": 0,
|
||||||
|
"zouit": 0,
|
||||||
|
"skipped": 0,
|
||||||
|
}
|
||||||
|
# grid-walk слоя CAT_PARCEL вернул layer_failed=True
|
||||||
|
mock_walk.return_value = (0, 9, True)
|
||||||
|
|
||||||
|
result = await harvest_quarter(
|
||||||
|
db=db,
|
||||||
|
client=client,
|
||||||
|
quarter="66:41:0303161",
|
||||||
|
job_id=1,
|
||||||
|
update_progress=lambda s: progress_states.append(s),
|
||||||
|
)
|
||||||
|
|
||||||
|
# failed_layers содержит CAT_PARCEL
|
||||||
|
assert result.failed_layers == [CAT_PARCEL]
|
||||||
|
# Финальный done-state несёт harvest_meta-флаг
|
||||||
|
assert result.phase_state is not None
|
||||||
|
assert result.phase_state["harvest_meta"] == {f"layer_{CAT_PARCEL}_failed": True}
|
||||||
|
# Квартал НЕ свалился — дошёл до done
|
||||||
|
assert progress_states[-1]["phase"] == "done"
|
||||||
|
|
||||||
|
|
||||||
|
# ── Issue #200: geom backfill для участков с geom IS NULL ─────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_mark_geom_unavailable_parametrized_no_fstring() -> None:
|
||||||
|
"""_mark_geom_unavailable использует ANY(CAST(:cads AS text[])) — не f-string SQL."""
|
||||||
|
from app.services.cadastre.bulk_harvest import _mark_geom_unavailable
|
||||||
|
|
||||||
|
executed: list[tuple[str, Any]] = []
|
||||||
|
|
||||||
|
db = MagicMock()
|
||||||
|
|
||||||
|
def capture(stmt: Any, params: Any = None) -> MagicMock:
|
||||||
|
executed.append((str(stmt), params))
|
||||||
|
return MagicMock(rowcount=2)
|
||||||
|
|
||||||
|
db.execute = capture
|
||||||
|
|
||||||
|
n = _mark_geom_unavailable(db, ["66:41:0303161:1", "66:41:0303161:2"])
|
||||||
|
|
||||||
|
assert n == 2
|
||||||
|
assert len(executed) == 1
|
||||||
|
sql, params = executed[0]
|
||||||
|
assert "ANY(CAST(:cads AS text[]))" in sql
|
||||||
|
assert "geom_unavailable = TRUE" in sql
|
||||||
|
assert params["cads"] == ["66:41:0303161:1", "66:41:0303161:2"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_mark_geom_unavailable_empty_noop() -> None:
|
||||||
|
"""Пустой список → 0, без db.execute."""
|
||||||
|
from app.services.cadastre.bulk_harvest import _mark_geom_unavailable
|
||||||
|
|
||||||
|
db = MagicMock()
|
||||||
|
assert _mark_geom_unavailable(db, []) == 0
|
||||||
|
db.execute.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_backfill_parcel_geom_heals_and_marks_unavailable() -> None:
|
||||||
|
"""Issue #200: участок с восстановленным geom → healed; оставшийся NULL → unavailable."""
|
||||||
|
from app.services.cadastre.bulk_harvest import backfill_parcel_geom
|
||||||
|
|
||||||
|
# Очередь: 2 участка одного квартала.
|
||||||
|
null_rows = ["66:41:0303161:1", "66:41:0303161:2"]
|
||||||
|
# После grid-walk участок :2 всё ещё NULL → должен быть помечен unavailable.
|
||||||
|
still_null = ["66:41:0303161:2"]
|
||||||
|
|
||||||
|
db = MagicMock()
|
||||||
|
captured: list[dict[str, Any]] = []
|
||||||
|
|
||||||
|
def execute_router(stmt: Any, params: Any = None) -> MagicMock:
|
||||||
|
sql = str(stmt)
|
||||||
|
result = MagicMock()
|
||||||
|
if "geom IS NULL AND geom_unavailable = FALSE" in sql and "LIMIT" in sql:
|
||||||
|
result.scalars.return_value.all.return_value = null_rows
|
||||||
|
elif "AND geom IS NULL" in sql and "ANY(CAST(:cads AS text[]))" in sql and "SELECT" in sql:
|
||||||
|
result.scalars.return_value.all.return_value = still_null
|
||||||
|
elif "UPDATE cad_parcels SET geom_unavailable = TRUE" in sql:
|
||||||
|
captured.append({"sql": sql, "params": params or {}})
|
||||||
|
result.rowcount = len(still_null)
|
||||||
|
else:
|
||||||
|
result.scalars.return_value.all.return_value = []
|
||||||
|
result.rowcount = 0
|
||||||
|
return result
|
||||||
|
|
||||||
|
db.execute = execute_router
|
||||||
|
db.commit = MagicMock()
|
||||||
|
|
||||||
|
client = AsyncMock()
|
||||||
|
|
||||||
|
with patch("app.services.cadastre.bulk_harvest._grid_walk_category") as mock_walk:
|
||||||
|
# grid-walk слоя 36368 «нашёл» 1 полигон, без layer_failed
|
||||||
|
mock_walk.return_value = (1, 9, False)
|
||||||
|
result = await backfill_parcel_geom(db=db, client=client, limit=500)
|
||||||
|
|
||||||
|
assert result.parcels_targeted == 2
|
||||||
|
assert result.quarters_scanned == 1
|
||||||
|
assert result.parcels_healed == 1 # :1 получил geom
|
||||||
|
assert result.parcels_marked_unavailable == 1 # :2 помечен
|
||||||
|
# grid-walk вызван для квартала с layer_id=36368
|
||||||
|
mock_walk.assert_called_once()
|
||||||
|
assert mock_walk.call_args.kwargs["quarter"] == "66:41:0303161"
|
||||||
|
assert mock_walk.call_args.kwargs["layer_id"] == 36368
|
||||||
|
# UPDATE geom_unavailable выполнен ровно для оставшегося NULL
|
||||||
|
assert len(captured) == 1
|
||||||
|
assert captured[0]["params"]["cads"] == still_null
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_backfill_parcel_geom_empty_queue_noop() -> None:
|
||||||
|
"""Issue #200: нет участков с geom IS NULL → пустой результат, grid-walk не вызывается."""
|
||||||
|
from app.services.cadastre.bulk_harvest import backfill_parcel_geom
|
||||||
|
|
||||||
|
db = MagicMock()
|
||||||
|
db.execute.return_value.scalars.return_value.all.return_value = []
|
||||||
|
client = AsyncMock()
|
||||||
|
|
||||||
|
with patch("app.services.cadastre.bulk_harvest._grid_walk_category") as mock_walk:
|
||||||
|
result = await backfill_parcel_geom(db=db, client=client, limit=500)
|
||||||
|
|
||||||
|
assert result.parcels_targeted == 0
|
||||||
|
assert result.quarters_scanned == 0
|
||||||
|
mock_walk.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_backfill_parcel_geom_marks_orphan_cads_unavailable() -> None:
|
||||||
|
"""Issue #200: cad без валидного 3-сегментного квартала → сразу unavailable (нет bbox)."""
|
||||||
|
from app.services.cadastre.bulk_harvest import backfill_parcel_geom
|
||||||
|
|
||||||
|
orphan_rows = ["66:41", "garbage"] # <3 сегментов → grid-walk невозможен
|
||||||
|
|
||||||
|
db = MagicMock()
|
||||||
|
marked: list[Any] = []
|
||||||
|
|
||||||
|
def execute_router(stmt: Any, params: Any = None) -> MagicMock:
|
||||||
|
sql = str(stmt)
|
||||||
|
result = MagicMock()
|
||||||
|
if "geom IS NULL AND geom_unavailable = FALSE" in sql and "LIMIT" in sql:
|
||||||
|
result.scalars.return_value.all.return_value = orphan_rows
|
||||||
|
elif "UPDATE cad_parcels SET geom_unavailable = TRUE" in sql:
|
||||||
|
marked.append(params)
|
||||||
|
result.rowcount = len(params["cads"])
|
||||||
|
else:
|
||||||
|
result.scalars.return_value.all.return_value = []
|
||||||
|
result.rowcount = 0
|
||||||
|
return result
|
||||||
|
|
||||||
|
db.execute = execute_router
|
||||||
|
db.commit = MagicMock()
|
||||||
|
client = AsyncMock()
|
||||||
|
|
||||||
|
with patch("app.services.cadastre.bulk_harvest._grid_walk_category") as mock_walk:
|
||||||
|
result = await backfill_parcel_geom(db=db, client=client, limit=500)
|
||||||
|
|
||||||
|
# Оба orphan помечены unavailable, grid-walk не вызывался
|
||||||
|
mock_walk.assert_not_called()
|
||||||
|
assert result.parcels_marked_unavailable == 2
|
||||||
|
assert marked and sorted(marked[0]["cads"]) == ["66:41", "garbage"]
|
||||||
|
|
|
||||||
51
data/sql/151_cad_parcels_geom_unavailable.sql
Normal file
51
data/sql/151_cad_parcels_geom_unavailable.sql
Normal file
|
|
@ -0,0 +1,51 @@
|
||||||
|
-- 151_cad_parcels_geom_unavailable.sql
|
||||||
|
-- ---------------------------------------------------------------------------
|
||||||
|
-- Контекст (issue #200):
|
||||||
|
-- На prod ~964 строки cad_parcels с geom IS NULL. Ground-truth (измерено
|
||||||
|
-- 2026-06-13): 930/964 — реальные участки ЕКБ 66:41:NNNNNN:NN, все
|
||||||
|
-- source='search'. NSPD search-эндпоинт для квартала возвращает участок без
|
||||||
|
-- полигона (только центроид/Point ИЛИ вовсе без geometry → upsert_parcel
|
||||||
|
-- пишет geom=NULL, см. bulk_harvest.py). Остаток (34) — кросс-региональные
|
||||||
|
-- артефакты поиска (последний сегмент :41 / :66 — это номер участка из
|
||||||
|
-- ДРУГОГО квартала, попавший в выдачу), у которых полигона в WMS-слое
|
||||||
|
-- 36368 нет вовсе.
|
||||||
|
--
|
||||||
|
-- Полигон для большинства восстановим grid-walk'ом слоя 36368 (WMS
|
||||||
|
-- GetFeatureInfo по bbox квартала — Phase 2-3 bulk_harvest). Но для
|
||||||
|
-- систематически недоступных (кросс-региональные артефакты) бесконечный
|
||||||
|
-- ретрай вреден: каждый прогон бэкфилла снова бы их дёргал.
|
||||||
|
--
|
||||||
|
-- Действие:
|
||||||
|
-- Добавить флаг geom_unavailable BOOLEAN. Backfill-task
|
||||||
|
-- (tasks.cadastre.backfill_parcel_geom) выставляет его в TRUE, когда после
|
||||||
|
-- grid-walk слоя участок ВСЁ ЕЩЁ без полигона — значит geom систематически
|
||||||
|
-- недоступен, перестаём ретраить. Дефолт FALSE — обычные участки и новые
|
||||||
|
-- inserts не затрагиваются; пересчёт делает только бэкфилл.
|
||||||
|
--
|
||||||
|
-- Partial-индекс под выборку «кого ещё докачивать»
|
||||||
|
-- (geom IS NULL AND NOT geom_unavailable) — узкий, индексирует только
|
||||||
|
-- рабочую очередь, не всю таблицу (40k строк).
|
||||||
|
--
|
||||||
|
-- Идемпотентность: ADD COLUMN IF NOT EXISTS + CREATE INDEX IF NOT EXISTS.
|
||||||
|
-- Apply order: чистая схема, без зависимости от кода. Backend-код с
|
||||||
|
-- backfill-task деплоится после (правило sql.md: схема → код).
|
||||||
|
-- ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
ALTER TABLE cad_parcels
|
||||||
|
ADD COLUMN IF NOT EXISTS geom_unavailable BOOLEAN NOT NULL DEFAULT FALSE;
|
||||||
|
|
||||||
|
COMMENT ON COLUMN cad_parcels.geom_unavailable IS
|
||||||
|
'issue #200: TRUE — geom систематически недоступен в НСПД (grid-walk слоя '
|
||||||
|
'36368 не вернул полигон). Выставляется backfill-task'
|
||||||
|
' (tasks.cadastre.backfill_parcel_geom), чтобы не ретраить докачку вечно. '
|
||||||
|
'FALSE (дефолт) — geom либо есть, либо ещё не пробовали докачать.';
|
||||||
|
|
||||||
|
-- Рабочая очередь бэкфилла: участки без геометрии, которые ещё имеет смысл
|
||||||
|
-- докачивать. Partial WHERE держит индекс компактным.
|
||||||
|
CREATE INDEX IF NOT EXISTS cad_parcels_geom_backfill_queue
|
||||||
|
ON cad_parcels (cad_num)
|
||||||
|
WHERE geom IS NULL AND geom_unavailable = FALSE;
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
Loading…
Add table
Reference in a new issue