50a261e05e
- 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/аудит/локальное чтение
358 lines
15 KiB
Python
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),
|
|
) |