ПТИЦА: убрать мёртвую фазу bulk_harvest, которая на каждый квартал запрашивала ПЗЗ у НСПД и писала их в таблицу без читателей #3579

Merged
bot-backend merged 4 commits from fix/ptica-bulk-harvest-phase25-removal into main 2026-09-17 11:12:02 +00:00
5 changed files with 46 additions and 352 deletions

View file

@ -587,19 +587,6 @@ class NSPDBulkClient:
)
return results
async def get_territorial_zones_in_bbox(
self,
bbox: tuple[float, float, float, float],
*,
grid_n: int = 7,
) -> list[dict]:
"""Grid-walk WMS GetFeatureInfo для layer 875838 (ПЗЗ территориальные зоны).
Returns: list of feature dicts с полями id, geometry, properties.
Дедуплицирует по feature id.
"""
return await self.get_features_in_bbox_grid(875838, bbox, grid_n=grid_n)
# ── 4. list_objects_in_building ───────────────────────────────────────────
# Q3 deferred — метод реализован, но не вызывается в bulk_harvest_quarter MVP.
# Готов для per-building помещения/парковка фазы.

View file

@ -19,7 +19,6 @@ Resumable: phase_state в cadastre_jobs показывает прогресс.
from __future__ import annotations
import hashlib
import json
import logging
from collections.abc import Callable
@ -260,10 +259,8 @@ async def harvest_quarter(
update_progress(done_progress)
# ── Phase 4: quarter stats + auto-heal geom из snapshot ─────────────────
# Bug #1583: auto-heal geom выполняем ДО Phase 2.5 (territorial_zones). Иначе
# кварталы с broken/NULL geom дают quarter_bbox_3857() == None → ПЗЗ молча
# пропускаются, а на следующем harvest квартал отсекается skip_fresh_hours.
# Чиним geom здесь → Phase 2.5 ниже получит валидный bbox в этом же прогоне.
# Bug #1583: кварталы с broken/NULL geom дают quarter_bbox_3857() == None →
# grid-walk следующего прогона молча пропускается. Чиним geom из snapshot.
stats_features = [f for f in snapshot.features if f.category_id == CAT_QUARTER_STATS]
if stats_features:
upsert_quarter_stats(db, quarter, stats_features[0])
@ -278,25 +275,6 @@ async def harvest_quarter(
logger.warning("harvest_quarter: geom auto-heal failed for %s: %s", quarter, e)
db.commit()
# ── Phase 2.5: grid-walk для territorial_zones (ПЗЗ, layer 875838) ────────
# Выполняем после основного grid-walk (Phase 2-3) И после Phase 4 geom
# auto-heal (см. Bug #1583) — так broken-geom кварталы, починенные выше,
# получают валидный bbox и ПЗЗ собираются в том же прогоне. Требует bbox квартала.
quarter_bbox = quarter_bbox_3857(db, quarter)
if quarter_bbox is not None:
update_progress({"phase": "territorial_zones_started", "quarter": quarter})
try:
tz_features = await client.get_territorial_zones_in_bbox(quarter_bbox)
tz_count = _save_territorial_zones(db, quarter, tz_features)
logger.info(
"harvest_quarter: territorial_zones quarter=%s upserted=%d", quarter, tz_count
)
except (NspdBulkWafError, NspdBulkRateLimitError):
# #2464-A: см. выше — бан пробрасываем, а не превращаем в «слой пуст».
raise
except Exception as e:
logger.warning("harvest_quarter: territorial_zones failed quarter=%s: %s", quarter, e)
# Issue #252: финальный phase_state несёт АГРЕГИРОВАННЫЙ harvest_meta по всем
# сбойным слоям. progress_cb мержит phase_state через JSONB `||` (shallow) —
# per-layer done-апдейты перетёрли бы harvest_meta друг друга, поэтому в
@ -679,7 +657,7 @@ def upsert_features(
cat = feature.category_id
# Per-feature SAVEPOINT (backend.md SAVEPOINT rule): один битый feature (edge-case
# GeoJSON / NOT NULL violation / type mismatch) НЕ должен ронять весь snapshot-tx
# квартала. Зеркалит _grid_walk_category (~368) и _save_territorial_zones (~1234):
# квартала. Зеркалит _grid_walk_category:
# begin_nested() rollback'ит только этот feature, остальные сохраняются. Счётчик
# инкрементим ТОЛЬКО при успехе (внутри блока) → counts остаются точными.
try:
@ -1476,114 +1454,6 @@ def upsert_quarter_stats(
)
# ── cad_territorial_zones upsert ────────────────────────────────────────────
def _save_territorial_zones(db: Session, quarter_cad: str, features: list[dict]) -> int:
"""UPSERT territorial_zones features в cad_territorial_zones по zone_id.
Args:
db: SQLAlchemy session.
quarter_cad: кадастровый номер квартала (3 сегмента).
features: list of raw feature dicts от get_features_in_bbox_grid.
Returns:
Количество успешно upserted строк.
"""
inserted = 0
for f in features:
props: dict = f.get("properties") or {}
geom = f.get("geometry")
geom_geojson: str | None = json.dumps(geom) if geom else None
# zone_id — NSPD feature id или стабильный fallback на основе md5 от properties.
# md5 гарантирует идемпотентность между runs (счётчик inserted сбрасывается).
_raw_id = props.get("id") or props.get("zone_id") or f.get("id")
if _raw_id:
zone_id = str(_raw_id)
else:
_props_hash = hashlib.md5(
json.dumps(props, sort_keys=True).encode("utf-8")
).hexdigest()[:12]
zone_id = f"{quarter_cad}_{_props_hash}"
zone_code = (
props.get("zone_code") or props.get("zone_index") or props.get("reg_numb_border")
)
zone_name = props.get("zone_name") or props.get("zone_type_name") or props.get("type_zone")
permitted_use = props.get("permitted_use") or props.get("vri")
# cad_territorial_zones.geom — geography(MultiPolygon, 4326).
# Polygon допустим (ST_Multi обернёт в SQL), но Point/LineString →
# geography-INSERT fail → SAVEPOINT откат → строка дропается молча.
# Зеркалит фильтр upsert_zouit (~1196).
geom_type = geom.get("type") if isinstance(geom, dict) else None
if geom_type not in ("Polygon", "MultiPolygon"):
if geom_type:
logger.info(
"_save_territorial_zones: zone_id=%s geom type=%s (не Polygon/MultiPolygon)"
" — geom=NULL",
zone_id,
geom_type,
)
geom_geojson = None
try:
# begin_nested() требует активной outer-транзакции для SAVEPOINT.
# SQLAlchemy Session (autobegin=True) автоматически начинает tx при первом
# db.execute() в этом loop — outer tx гарантирована.
with db.begin_nested():
db.execute(
text("""
INSERT INTO cad_territorial_zones
(quarter_cad, zone_id, zone_code, zone_name,
permitted_use, raw_props, geom)
VALUES (
CAST(:quarter_cad AS text),
CAST(:zone_id AS text),
CAST(:zone_code AS text),
CAST(:zone_name AS text),
CAST(:permitted_use AS text),
CAST(:raw_props AS jsonb),
CASE WHEN CAST(:geom AS text) IS NOT NULL
THEN ST_Multi(
ST_Transform(
ST_SetSRID(ST_GeomFromGeoJSON(CAST(:geom AS text)), 3857),
4326
)
)::geography
ELSE NULL
END
)
ON CONFLICT (zone_id) DO UPDATE SET
zone_code = EXCLUDED.zone_code,
zone_name = EXCLUDED.zone_name,
permitted_use = EXCLUDED.permitted_use,
raw_props = EXCLUDED.raw_props,
geom = EXCLUDED.geom,
fetched_at = NOW()
"""),
{
"quarter_cad": quarter_cad,
"zone_id": zone_id,
"zone_code": zone_code,
"zone_name": zone_name,
"permitted_use": permitted_use,
"raw_props": json.dumps(props, ensure_ascii=False),
"geom": geom_geojson,
},
)
inserted += 1
except Exception as e:
logger.warning(
"_save_territorial_zones: upsert failed zone_id=%s quarter=%s: %s",
zone_id,
quarter_cad,
e,
)
db.commit()
return inserted
# ── Утилиты ──────────────────────────────────────────────────────────────────

View file

@ -329,7 +329,7 @@ def build_beat_schedule() -> dict:
# вывел PKK6-эндпоинт), таблица pzz_zones_ekb пуста (0 rows на prod) — задача ни
# разу не наполнила её успешно. ПЗЗ-данные пришли в систему ИНЫМ путём:
# zone_regulation_cache (#1059, beat zone-regulation-refresh-monthly) + NSPD
# territorial_zones dumps (Phase 2.5 bulk_harvest, layer 875838). Активный beat
# territorial_zones в nspd_quarter_dumps (nspd_sync, layer 875838). Активный beat
# каждый месяц дёргал deprecated PKK6 → broad except в pzz_sync.py логировал
# ошибку как error → рекуррентил GlitchTip BACKEND-1B.
#

View file

@ -1,198 +0,0 @@
"""Тесты для _save_territorial_zones (bulk_harvest.py) — mock-based.
Проверяет:
- Успешный UPSERT 3 features 3 строки вставлены
- Повторный вызов ON CONFLICT обновляет, не дублирует
- Feature без geometry строка вставлена с geom=NULL, без краша
- Feature без zone_id синтетический fallback zone_id используется
"""
from __future__ import annotations
from typing import Any
from unittest.mock import MagicMock
from app.services.cadastre.bulk_harvest import _save_territorial_zones
def _make_feature(
feature_id: Any = "zone_1",
zone_code: str = "Ж-1",
zone_name: str = "Жилая смешанная",
permitted_use: str = "ИЖС",
has_geometry: bool = True,
) -> dict:
"""Создать raw feature dict в формате get_features_in_bbox_grid."""
geom = (
{
"type": "Polygon",
"coordinates": [
[
[6090000.0, 7590000.0],
[6090100.0, 7590000.0],
[6090100.0, 7590100.0],
[6090000.0, 7590100.0],
[6090000.0, 7590000.0],
]
],
}
if has_geometry
else None
)
return {
"id": feature_id,
"geometry": geom,
"properties": {
"zone_code": zone_code,
"zone_name": zone_name,
"permitted_use": permitted_use,
},
}
def _make_db_mock() -> MagicMock:
"""Mock SQLAlchemy Session с begin_nested() savepoint support."""
db = MagicMock()
# begin_nested() используется как context manager
savepoint_ctx = MagicMock()
savepoint_ctx.__enter__ = MagicMock(return_value=savepoint_ctx)
savepoint_ctx.__exit__ = MagicMock(return_value=False)
db.begin_nested.return_value = savepoint_ctx
return db
class TestSaveTerritorialZones:
"""Тесты для _save_territorial_zones."""
def test_three_features_inserted(self) -> None:
"""3 features → returned count == 3, execute вызван 3 раза."""
db = _make_db_mock()
features = [
_make_feature("z1", "Ж-1"),
_make_feature("z2", "ОД-1"),
_make_feature("z3", "П-1"),
]
result = _save_territorial_zones(db, "66:41:0204016", features)
assert result == 3
assert db.execute.call_count == 3
db.commit.assert_called_once()
def test_empty_features_list(self) -> None:
"""Пустой список → 0 inserted, commit всё равно вызван."""
db = _make_db_mock()
result = _save_territorial_zones(db, "66:41:0204016", [])
assert result == 0
db.execute.assert_not_called()
db.commit.assert_called_once()
def test_feature_without_geometry_no_crash(self) -> None:
"""Feature без geometry → geom=NULL, строка вставлена без краша."""
db = _make_db_mock()
features = [_make_feature("zone_no_geom", has_geometry=False)]
result = _save_territorial_zones(db, "66:41:0204016", features)
assert result == 1
# Проверяем что geom параметр передан как None
call_kwargs: dict = db.execute.call_args[0][1]
assert call_kwargs["geom"] is None
def test_feature_without_zone_id_uses_fallback(self) -> None:
"""Feature без id → md5-based fallback zone_id (stable между runs)."""
db = _make_db_mock()
features = [
{
"id": None,
"geometry": None,
"properties": {"zone_code": "Ж-2"},
}
]
result = _save_territorial_zones(db, "66:41:0204016", features)
assert result == 1
call_kwargs = db.execute.call_args[0][1]
zone_id: str = call_kwargs["zone_id"]
# fallback zone_id содержит quarter_cad и стабильный hash (12 hex chars)
assert zone_id.startswith("66:41:0204016_")
suffix = zone_id.split("_", 3)[-1]
assert len(suffix) == 12
assert all(c in "0123456789abcdef" for c in suffix)
# Второй вызов с теми же данными → тот же zone_id (идемпотентность)
db2 = _make_db_mock()
_save_territorial_zones(db2, "66:41:0204016", features)
call_kwargs2 = db2.execute.call_args[0][1]
assert call_kwargs2["zone_id"] == zone_id
def test_zone_id_from_props_id(self) -> None:
"""Если feature.id=None, но props['id'] есть — используется props['id']."""
db = _make_db_mock()
features = [
{
"id": None,
"geometry": None,
"properties": {"id": "props_id_42", "zone_code": "Ж-3"},
}
]
result = _save_territorial_zones(db, "66:41:0204016", features)
assert result == 1
call_kwargs = db.execute.call_args[0][1]
assert call_kwargs["zone_id"] == "props_id_42"
def test_execute_error_logged_not_raised(self) -> None:
"""Exception в execute → строка не вставлена, warning залогирован, не re-raise."""
db = _make_db_mock()
db.execute.side_effect = RuntimeError("DB error")
features = [_make_feature("z_err")]
# Не должен бросить исключение
result = _save_territorial_zones(db, "66:41:0204016", features)
assert result == 0
db.commit.assert_called_once()
def test_savepoint_used_per_row(self) -> None:
"""begin_nested() вызывается для каждой строки (SAVEPOINT паттерн)."""
db = _make_db_mock()
features = [_make_feature(f"z{i}") for i in range(3)]
_save_territorial_zones(db, "66:41:0204016", features)
assert db.begin_nested.call_count == 3
def test_quarter_cad_param_passed(self) -> None:
"""quarter_cad правильно передаётся в SQL параметры."""
db = _make_db_mock()
features = [_make_feature("zone_check")]
_save_territorial_zones(db, "66:41:9999999", features)
call_kwargs = db.execute.call_args[0][1]
assert call_kwargs["quarter_cad"] == "66:41:9999999"
def test_raw_props_serialized(self) -> None:
"""raw_props — JSON строка из properties dict."""
import json
db = _make_db_mock()
features = [
{
"id": "z_props",
"geometry": None,
"properties": {"zone_code": "ОД-2", "extra": "value"},
}
]
_save_territorial_zones(db, "66:41:0204016", features)
call_kwargs = db.execute.call_args[0][1]
raw = json.loads(call_kwargs["raw_props"])
assert raw["zone_code"] == "ОД-2"
assert raw["extra"] == "value"

View file

@ -489,7 +489,7 @@ async def test_harvest_quarter_does_not_early_exit_on_shared_phase_done() -> Non
db = MagicMock()
# Симулируем shared phase_state с phase=done от ДРУГОГО quarter.
# xmin=None → quarter_bbox_3857 (grid-walk geometry helper) вернёт None →
# grid-walk + territorial_zones фазы корректно пропускаются (тест про
# grid-walk фаза корректно пропускается (тест про
# snapshot/idempotency, не про geometry). Тот же dict возвращается на ВСЕ
# db.execute().mappings().first() в этом тесте.
db.execute = MagicMock(
@ -554,7 +554,7 @@ async def test_harvest_quarter_calls_upsert_features() -> None:
db = MagicMock()
# phase_state = None → начинаем с нуля
# xmin=None → quarter_bbox_3857 (grid-walk geometry helper) вернёт None →
# grid-walk + territorial_zones фазы пропускаются (тесты про snapshot /
# grid-walk фаза пропускается (тесты про snapshot /
# per-cat-probe, не про geometry). Тот же dict на ВСЕ
# db.execute().mappings().first() вызовы.
db.execute = MagicMock(
@ -759,7 +759,7 @@ async def test_harvest_quarter_calls_per_cat_probe_for_zouit_when_meta_nonzero()
db = MagicMock()
# xmin=None → quarter_bbox_3857 (grid-walk geometry helper) вернёт None →
# grid-walk + territorial_zones фазы пропускаются (тесты про snapshot /
# grid-walk фаза пропускается (тесты про snapshot /
# per-cat-probe, не про geometry). Тот же dict на ВСЕ
# db.execute().mappings().first() вызовы.
db.execute = MagicMock(
@ -834,7 +834,7 @@ async def test_harvest_quarter_skips_per_cat_probe_when_meta_zero() -> None:
db = MagicMock()
# xmin=None → quarter_bbox_3857 (grid-walk geometry helper) вернёт None →
# grid-walk + territorial_zones фазы пропускаются (тесты про snapshot /
# grid-walk фаза пропускается (тесты про snapshot /
# per-cat-probe, не про geometry). Тот же dict на ВСЕ
# db.execute().mappings().first() вызовы.
db.execute = MagicMock(
@ -904,7 +904,7 @@ async def test_harvest_quarter_per_cat_probe_enk_called_when_meta_nonzero() -> N
db = MagicMock()
# xmin=None → quarter_bbox_3857 (grid-walk geometry helper) вернёт None →
# grid-walk + territorial_zones фазы пропускаются (тесты про snapshot /
# grid-walk фаза пропускается (тесты про snapshot /
# per-cat-probe, не про geometry). Тот же dict на ВСЕ
# db.execute().mappings().first() вызовы.
db.execute = MagicMock(
@ -1223,7 +1223,7 @@ async def test_harvest_quarter_geom_heal_failure_does_not_propagate() -> None:
db = MagicMock()
# xmin=None → quarter_bbox_3857 (grid-walk geometry helper) вернёт None →
# grid-walk + territorial_zones фазы пропускаются (тесты про snapshot /
# grid-walk фаза пропускается (тесты про snapshot /
# per-cat-probe, не про geometry). Тот же dict на ВСЕ
# db.execute().mappings().first() вызовы.
db.execute = MagicMock(
@ -1518,7 +1518,6 @@ async def test_harvest_quarter_records_failed_layer_in_phase_state() -> None:
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]] = []
@ -1557,6 +1556,42 @@ async def test_harvest_quarter_records_failed_layer_in_phase_state() -> None:
assert progress_states[-1]["phase"] == "done"
@pytest.mark.asyncio
async def test_harvest_quarter_makes_no_territorial_zones_request_2985() -> None:
"""#2985: Phase 2.5 удалена. Квартал без overflow с валидным bbox стоит ровно один
запрос к НСПД (search_by_quarter) отдельного grid-walk за ПЗЗ больше нет."""
from app.services.cadastre.bulk_harvest import harvest_quarter
snapshot = QuarterSnapshot(
quarter="66:41:0303161",
fetched_at="2026-05-15T10:00:00+00:00",
features=[_make_parcel_feature()],
meta_counts={},
)
client = AsyncMock()
client.search_by_quarter = AsyncMock(return_value=snapshot)
progress_states: list[dict[str, Any]] = []
with (
patch("app.services.cadastre.bulk_harvest.upsert_features") as mock_upsert,
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 = dict.fromkeys(
("parcels", "buildings", "constructions", "oncs", "enks", "zouit", "skipped"), 0
)
await harvest_quarter(
db=_mock_db_grid_bbox(),
client=client,
quarter="66:41:0303161",
job_id=1,
update_progress=progress_states.append,
)
assert [c[0] for c in client.mock_calls] == ["search_by_quarter"]
assert [s["phase"] for s in progress_states] == ["snapshot_started", "snapshot_done", "done"]
# ── Issue #200: geom backfill для участков с geom IS NULL ─────────────────────