Files
Sadmin 50a261e05e B17/E5: ингестия полевых данных КОНТУРа
- POST /operations/{id}/field-data/sync: тянет GET /api/v1/vector/areas-checked
  → идемпотентный upsert в areas_checked (dedupe по contour_area_id:
  geom/checked_at/result/coverage_pct/team_label, source=contour)
- GET /operations/{id}/areas-checked: теперь из локальной БД после sync
  (source=db), фолбэк на живой прокси КОНТУРа если пусто (source=contour/empty)
- миграция 011_b17_field_data: areas_checked +contour_area_id/team_label/source
  (+check ck_areas_checked_source)
- аудит field_data_sync; лог vector.contour field_sync_*
- фронт: кнопка «⟳ Синхронизировать полевые данные» в КОНТУР-панели,
  слой проверенных квадратов на карте (зелёный=clear ≥99.5%, жёлтый=partial)
- тесты: +9 (267 passed) — sync/идемпотентность/502/аудит/локальное чтение
2026-09-25 10:34:53 +03:00

358 lines
15 KiB
Python

"""B21/E4: отправка приоритетных зон в КОНТУР (контракт B20, outbound).
POST /api/v1/operations/{id}/send-to-contour:
анализ кейса → SearchZones[] (GeoJSON-сектора направлений) →
POST {CONTOUR_API_URL}/api/v1/vector/zones?operation_id=<contour_operation_id>
B17/E5: ингестия полевых данных (inbound):
POST /api/v1/operations/{id}/field-data/sync — тянет из КОНТУРа
GET /api/v1/vector/areas-checked?operation_id=<contour_operation_id> и
upsert'ит в локальную areas_checked (dedupe по contour_area_id);
GET /api/v1/operations/{id}/areas-checked — читает локальную БД.
Требуется: у операции задан contour_operation_id (UUID операции в КОНТУРе),
env CONTOUR_API_URL и CONTOUR_TOKEN (межсервисный, в secrets).
"""
from __future__ import annotations
import json
import math
import os
import uuid as uuid_mod
from datetime import datetime, timezone
from typing import Any, Optional
import httpx
from fastapi import APIRouter, Depends, HTTPException, Request
from pydantic import BaseModel
from sqlalchemy.orm import Session
from backend.audit import audit_log, can_access_unit
from backend.database import get_db
from backend.logging_setup import get_logger
from backend.models import AreaChecked, SearchOperation, User
from backend.routers.auth import get_current_user, require_permission
log = get_logger('contour')
router = APIRouter(prefix='/api/v1/operations', tags=['operations-contour'])
def _contour_url() -> str:
return os.getenv('CONTOUR_API_URL', 'http://192.168.0.130:8000')
def _contour_token() -> str:
return os.getenv('CONTOUR_TOKEN', '')
SCHEMA_VERSION = '1.0'
_DIRECTIONS = {
'N': 0, 'NE': 45, 'E': 90, 'SE': 135,
'S': 180, 'SW': 225, 'W': 270, 'NW': 315,
}
class SendResult(BaseModel):
sent_zones: int
contour_operation_id: str
contour_response: dict[str, Any]
def _sector_polygon(lat: float, lon: float, radius_km: float, direction: str,
segments: int = 24) -> dict:
"""GeoJSON Polygon сектора ±22.5° вокруг направления (для контракта B20)."""
az = _DIRECTIONS.get(direction, 0)
half = 22.5
ring: list[list[float]] = [[lon, lat]]
for i in range(segments + 1):
bearing = math.radians(az - half + (2 * half) * i / segments)
dlat = (radius_km / 111.0) * math.cos(bearing)
dlon = (radius_km / (111.0 * math.cos(math.radians(lat)))) * math.sin(bearing)
ring.append([round(lon + dlon, 7), round(lat + dlat, 7)])
ring.append([lon, lat])
return {'type': 'Polygon', 'coordinates': [ring]}
def _probability_from_priority(priority: int, total: int) -> float:
"""Градация: топ-зона → наибольшая вероятность (0..1)."""
if total <= 1:
return 0.8
return round(max(0.1, 0.9 - 0.7 * (priority - 1) / max(1, total - 1)), 3)
def _load_analysis(db, case_id) -> dict:
"""Последний анализ кейса: search_models (B15) или analysis_log (legacy)."""
from sqlalchemy import text as sa_text
row = db.execute(
sa_text("SELECT model_json FROM search_models "
"WHERE case_id = :cid ORDER BY version DESC LIMIT 1"),
{'cid': str(case_id)},
).first()
if row and row[0]:
return dict(row[0])
from backend.models import Case
case = db.get(__import__('backend.models', fromlist=['Case']).Case, case_id)
if case and case.analysis_log and isinstance(case.analysis_log, dict):
return dict(case.analysis_log)
raise HTTPException(status_code=400, detail='У карточки нет анализа — запустите анализ сначала')
def _build_b20_payload(op: SearchOperation, analysis: dict) -> dict:
"""SearchZones[] из primary_zones анализа (полигоны направлений)."""
lat = analysis.get('lat') or analysis.get('tnp_lat')
lon = analysis.get('lon') or analysis.get('tnp_lon')
if lat is None or lon is None:
raise HTTPException(status_code=400,
detail='В анализе нет координат ТНП — невозможны геометрии зон')
max_distance = float(analysis.get('max_distance_km') or 3.0)
zones_in = analysis.get('primary_zones') or []
if not zones_in:
raise HTTPException(status_code=400, detail='В анализе нет приоритетных зон')
out_zones = []
for z in zones_in[:5]:
direction = z.get('direction')
distance = float(z.get('distance') or max_distance * 0.5)
priority = int(z.get('priority') or 1)
out_zones.append({
'zone_id': f"vec:{direction}:{int(distance * 1000)}m",
'geom': _sector_polygon(float(lat), float(lon), distance, direction),
'priority': priority,
'probability': _probability_from_priority(priority, len(zones_in[:5])),
'reasoning': z.get('reason') or f'Приоритет {priority} по модели ВЕКТОРа',
})
return {
'schema_version': SCHEMA_VERSION,
'case_id': str(op.case_id),
'model_version': f"vector:{(analysis.get('analyzed_at') or datetime.now(timezone.utc).isoformat())[:19]}",
'generated_at': analysis.get('analyzed_at') or datetime.now(timezone.utc).isoformat(),
'max_distance_km': max_distance,
'zones': out_zones,
}
@router.post('/{operation_id}/send-to-contour')
def send_to_contour(
operation_id: str,
request: Request,
current_user: User = Depends(require_permission('update')),
db = Depends(get_db),
) -> SendResult:
op = _load_operation(db, current_user, operation_id)
if not op.contour_operation_id:
raise HTTPException(
status_code=400,
detail='У операции не задан contour_operation_id (UUID операции в КОНТУРе) — '
'укажите его в карточке операции (PATCH /operations/{id})',
)
if not _contour_token():
raise HTTPException(status_code=503, detail='CONTOUR_TOKEN не настроен')
analysis = _load_analysis(db, op.case_id)
payload = _build_b20_payload(op, analysis)
url = f"{_contour_url()}/api/v1/vector/zones?operation_id={op.contour_operation_id}"
try:
resp = httpx.post(url, json=payload,
headers={'Authorization': f'Bearer {_contour_token()}'},
timeout=30.0)
except httpx.HTTPError as e:
log.warning('contour_send_failed op=%s url=%s error=%s', op.id, _contour_url(), e)
audit_log(db, current_user, 'contour_send_failed', object_type='operation',
object_id=str(op.id), request=request, details={'error': str(e)})
raise HTTPException(status_code=502, detail=f'КОНТУР недоступен: {e}') from e
if resp.status_code >= 400:
log.warning('contour_rejected op=%s status=%d body=%s',
op.id, resp.status_code, resp.text[:200])
audit_log(db, current_user, 'contour_send_failed', object_type='operation',
object_id=str(op.id), request=request,
details={'status': resp.status_code, 'body': resp.text[:500]})
raise HTTPException(status_code=502,
detail=f'КОНТУР отклонил импорт: {resp.status_code} {resp.text[:300]}')
result = resp.json()
log.info('contour_send_ok op=%s зон=%d contour_op=%s',
op.id, len(payload['zones']), op.contour_operation_id)
audit_log(db, current_user, 'contour_send_ok', object_type='operation',
object_id=str(op.id), request=request,
details={'zones': result.get('zones_total'), 'status': result.get('status')})
return SendResult(
sent_zones=payload['zones'].__len__(),
contour_operation_id=str(op.contour_operation_id),
contour_response=result,
)
def _load_operation(db, user: User, operation_id: str) -> SearchOperation:
op = db.get(SearchOperation, uuid_mod.UUID(operation_id)) if operation_id else None
if not op:
raise HTTPException(status_code=404, detail='Операция не найдена')
if not can_access_unit(db, user, str(op.unit_id) if op.unit_id else None):
raise HTTPException(status_code=403, detail='Операция вне вашего скоупа')
return op
@router.get('/{operation_id}/areas-checked')
def areas_checked(
operation_id: str,
current_user: User = Depends(get_current_user),
db: Session = Depends(get_db),
) -> dict[str, Any]:
"""Проверенные квадраты операции (B17): из локальной БД после sync.
Если ещё не синхронизировано (0 записей) — прозрачно проксирует в КОНТУР,
чтобы страница работала без обязательного sync (live-режим, B17-заготовка).
"""
op = _load_operation(db, current_user, operation_id)
local = _local_areas_checked(db, op)
if local:
return {'source': 'db', 'items': local}
# Фолбэк: живой прокси в КОНТУР (как было до B17)
if not op.contour_operation_id:
return {'source': 'empty', 'items': []}
if not _contour_token():
raise HTTPException(status_code=503, detail='CONTOUR_TOKEN не настроен')
url = f"{_contour_url()}/api/v1/vector/areas-checked?operation_id={op.contour_operation_id}"
try:
resp = httpx.get(url, headers={'Authorization': f'Bearer {_contour_token()}'}, timeout=30.0)
except httpx.HTTPError as e:
raise HTTPException(status_code=502, detail=f'КОНТУР недоступен: {e}') from e
if resp.status_code >= 400:
raise HTTPException(status_code=502, detail=f'КОНТУР: {resp.status_code} {resp.text[:300]}')
return {'source': 'contour', 'items': resp.json()}
def _parse_iso(dt_str: Any) -> Optional[datetime]:
if not dt_str or not isinstance(dt_str, str):
return None
try:
return datetime.fromisoformat(str(dt_str).replace('Z', '+00:00'))
except ValueError:
return None
def _parse_contour_geom(raw: Any) -> Optional[dict]:
"""Геометрия КОНТУРа приходит JSON-строкой (ST_AsGeoJSON) — иногда dict."""
if raw is None:
return None
if isinstance(raw, dict):
return raw
try:
return json.loads(str(raw))
except (ValueError, TypeError):
return None
def _local_areas_checked(db: Session, op: SearchOperation) -> list[dict[str, Any]]:
"""Локальные проверенные квадраты операции (после field-data/sync)."""
from backend.models import Case
case = db.get(Case, op.case_id)
if case is None:
return []
rows = (db.query(AreaChecked)
.filter(AreaChecked.case_id == op.case_id)
.order_by(AreaChecked.checked_at.desc().nullslast())
.all())
return [{
'area_id': r.contour_area_id or str(r.id),
'geom': r.geom,
'team_id': r.team_label,
'checked_at': r.checked_at.isoformat() if r.checked_at else None,
'result': r.result,
'coverage_pct': r.coverage_pct,
'synced': True,
} for r in rows]
class SyncResult(BaseModel):
synced_areas: int
updated_areas: int
new_areas: int
found_events: int
last_checked_at: Optional[str] = None
contour_area_count: int
@router.post('/{operation_id}/field-data/sync')
def field_data_sync(
operation_id: str,
request: Request,
current_user: User = Depends(require_permission('update')),
db: Session = Depends(get_db),
) -> SyncResult:
"""B17: синхронизация полевых данных из КОНТУРа в локальную БД.
Тянет GET /api/v1/vector/areas-checked (весь список, idempotent upsert
по contour_area_id). found (100%) на стороне КОНТУРа сейчас не бывает —
находка фиксируется завершением поиска; счётчик found_events — резерв.
"""
op = _load_operation(db, current_user, operation_id)
if not op.contour_operation_id:
raise HTTPException(
status_code=400,
detail='У операции не задан contour_operation_id — укажите его в карточке операции',
)
if not _contour_token():
raise HTTPException(status_code=503, detail='CONTOUR_TOKEN не настроен')
url = f"{_contour_url()}/api/v1/vector/areas-checked?operation_id={op.contour_operation_id}"
try:
resp = httpx.get(url, headers={'Authorization': f'Bearer {_contour_token()}'}, timeout=30.0)
except httpx.HTTPError as e:
log.warning('field_sync_unreachable op=%s error=%s', op.id, e)
raise HTTPException(status_code=502, detail=f'КОНТУР недоступен: {e}') from e
if resp.status_code >= 400:
log.warning('field_sync_rejected op=%s status=%d', op.id, resp.status_code)
raise HTTPException(status_code=502, detail=f'КОНТУР: {resp.status_code} {resp.text[:300]}')
items = resp.json()
if not isinstance(items, list):
raise HTTPException(status_code=502, detail='КОНТУР вернул неожиданный формат areas-checked')
case_id = op.case_id
existing = {r.contour_area_id: r for r in
db.query(AreaChecked).filter(AreaChecked.case_id == case_id).all()
if r.contour_area_id}
new_n = upd_n = 0
latest: Optional[datetime] = None
for item in items:
if not isinstance(item, dict) or not item.get('area_id'):
continue
geom = _parse_contour_geom(item.get('geom'))
checked_at = _parse_iso(item.get('checked_at'))
if checked_at and (latest is None or checked_at > latest):
latest = checked_at
row = existing.get(item['area_id'])
if row is None:
row = AreaChecked(case_id=case_id, contour_area_id=item['area_id'])
db.add(row)
new_n += 1
else:
upd_n += 1
row.geom = geom
row.checked_at = checked_at
row.result = 'clear' if item.get('result') == 'clear' else 'partial'
row.coverage_pct = float(item['coverage_pct']) if item.get('coverage_pct') is not None else None
row.team_label = str(item['team_id']) if item.get('team_id') else None
row.source = 'contour'
db.commit()
if new_n or upd_n:
log.info('field_sync op=%s квадратов=%d (новых=%d обновлено=%d) from=%s',
op.id, new_n + upd_n, new_n, _contour_url())
audit_log(db, current_user, 'field_data_sync', object_type='search_operation',
object_id=str(op.id), request=request,
details={'new': new_n, 'updated': upd_n,
'total': len(items)})
return SyncResult(
synced_areas=new_n + upd_n,
updated_areas=upd_n,
new_areas=new_n,
found_events=0,
last_checked_at=latest.isoformat() if latest else None,
contour_area_count=len(items),
)