B2: tolerant JSON extraction from Claude responses + graceful fallback

_extract_json_payload handles a json/JSON/bare fence, raw JSON and JSON
embedded in prose; any unparseable response or contract violation now
degrades to the deterministic scoring_service instead of raising.
Also guards the response envelope itself (content[0].text).

B3: single home for recommendation scoring

services/recommendation_service.py holds the rules; routers/stats.py and
backend/services/stats_service.py both delegate to it. Unified rules are the
union of the two old copies: same weights/threshold, substring matching
(superset of the old exact match), tolerant key aliases, health_flags rule
kept. Endpoint response contract unchanged.

Plus: Overpass circuit breaker and concurrent zone queries in geo_service -
128 sequential calls per analysis no longer each burn a connect timeout when
the host has no outbound network.

Tests: 152 -> 194 passed.
This commit is contained in:
root
2026-07-25 13:01:02 +00:00
parent 76278f5fe1
commit 8b3a2cbf7e
7 changed files with 509 additions and 59 deletions
+47 -12
View File
@@ -171,21 +171,56 @@ async def analyze_with_claude(case_data: dict, api_key: str) -> AnalysisResult:
if response.status_code != 200:
raise Exception(f"Anthropic API error: {response.status_code} - {response.text}")
result = response.json()
content = result["content"][0]["text"]
try:
result = response.json()
content = result["content"][0]["text"]
analysis_data = _extract_json_payload(content)
analysis_data['fallback_used'] = False
return AnalysisResult(**analysis_data)
except Exception as e:
# Ответ модели пришёл в неожидаемом виде — не роняем анализ,
# а отдаём детерминированный результат scoring_service.
logger.warning(
f"Не удалось разобрать ответ Claude ({type(e).__name__}: {e}). "
"Используется fallback scoring service."
)
return await analyze_with_fallback(case_data)
# Парсим JSON из ответа
# Убираем возможные markdown блоки кода
if "```json" in content:
content = content.split("```json")[1].split("```")[0].strip()
elif "```" in content:
content = content.split("```")[1].split("```")[0].strip()
analysis_data = json.loads(content)
analysis_data['fallback_used'] = False
def _extract_json_payload(content: str) -> dict:
"""Extract a JSON object from a model response.
# Преобразуем в Pydantic модель
return AnalysisResult(**analysis_data)
Tolerates a ```json fence, a bare ``` fence, or raw JSON with
surrounding prose. Raises ValueError if nothing parseable is found.
"""
candidates = []
stripped = (content or '').strip()
if '```' in stripped:
for marker in ('```json', '```JSON', '```'):
if marker in stripped:
after = stripped.split(marker, 1)[1]
candidates.append(after.split('```', 1)[0].strip())
break
candidates.append(stripped)
# Last resort: the widest {...} span in the text.
start, end = stripped.find('{'), stripped.rfind('}')
if start != -1 and end > start:
candidates.append(stripped[start:end + 1])
for candidate in candidates:
if not candidate:
continue
try:
parsed = json.loads(candidate)
except (json.JSONDecodeError, TypeError):
continue
if isinstance(parsed, dict):
return parsed
raise ValueError('Не удалось извлечь JSON из ответа модели')
async def analyze_with_fallback(case_data: dict) -> AnalysisResult:
+72 -8
View File
@@ -1,9 +1,12 @@
"""
Geo service for building search zones and querying OpenStreetMap data via Overpass API.
"""
import asyncio
import math
import json
import hashlib
import logging
import time
from datetime import datetime, timedelta
from pathlib import Path
from typing import List, Dict, Optional, Tuple
@@ -26,6 +29,52 @@ CACHE_DIR = Path("/tmp/overpass_cache")
CACHE_TTL_HOURS = 24
OVERPASS_URL = "https://overpass-api.de/api/interpreter"
logger = logging.getLogger(__name__)
# --- Overpass circuit breaker ---------------------------------------------
# build_search_zones issues ~128 Overpass calls per analysis. When the host has
# no outbound connectivity every one of them burns the full connect timeout,
# which turns a single analysis into several minutes of waiting for results
# that are empty anyway. After OVERPASS_FAILURE_THRESHOLD consecutive
# transport failures we stop calling out until OVERPASS_COOLDOWN_SECONDS have
# passed. Callers get the same {'elements': []} they already got on error.
OVERPASS_FAILURE_THRESHOLD = 3
OVERPASS_COOLDOWN_SECONDS = 60.0
OVERPASS_CONNECT_TIMEOUT = 5.0
_overpass_failures = 0
_overpass_open_until = 0.0
def _overpass_circuit_open() -> bool:
"""True while the breaker is tripped (skip network, return empty fast)."""
if _overpass_failures < OVERPASS_FAILURE_THRESHOLD:
return False
if time.monotonic() >= _overpass_open_until:
_reset_overpass_circuit()
return False
return True
def _record_overpass_failure() -> None:
global _overpass_failures, _overpass_open_until
_overpass_failures += 1
if _overpass_failures == OVERPASS_FAILURE_THRESHOLD:
_overpass_open_until = time.monotonic() + OVERPASS_COOLDOWN_SECONDS
logger.warning(
"Overpass API недоступен (%d подряд неудачных запросов). "
"Геоданные отключены на %.0f c, анализ продолжается без них.",
_overpass_failures,
OVERPASS_COOLDOWN_SECONDS,
)
def _reset_overpass_circuit() -> None:
global _overpass_failures, _overpass_open_until
_overpass_failures = 0
_overpass_open_until = 0.0
# Direction mappings
SEARCH_DISTANCES = [500, 1000, 2000, 5000]
@@ -167,8 +216,13 @@ async def query_overpass(query: str) -> Dict:
if cached is not None:
return cached
# Skip the network entirely while the breaker is tripped.
if _overpass_circuit_open():
return {'elements': []}
# Query API
async with httpx.AsyncClient(timeout=30.0) as client:
timeout = httpx.Timeout(30.0, connect=OVERPASS_CONNECT_TIMEOUT)
async with httpx.AsyncClient(timeout=timeout) as client:
try:
response = await client.post(
OVERPASS_URL,
@@ -181,8 +235,10 @@ async def query_overpass(query: str) -> Dict:
# Save to cache
save_to_cache(cache_key, data)
_reset_overpass_circuit()
return data
except Exception as e:
_record_overpass_failure()
# Return empty result on error
return {'elements': []}
@@ -296,8 +352,6 @@ async def get_zone_features(lat: float, lon: float, direction: str, radius_m: in
);
out geom;
"""
roads_data = await query_overpass(roads_query)
roads_km = calculate_road_length(roads_data.get('elements', []))
# Query water bodies
water_query = f"""
@@ -308,8 +362,7 @@ async def get_zone_features(lat: float, lon: float, direction: str, radius_m: in
);
out center;
"""
water_data = await query_overpass(water_query)
water_distance = find_nearest_distance(lat, lon, water_data.get('elements', []))
# Query settlements
settlement_query = f"""
@@ -319,8 +372,7 @@ async def get_zone_features(lat: float, lon: float, direction: str, radius_m: in
);
out;
"""
settlement_data = await query_overpass(settlement_query)
settlement_distance = find_nearest_distance(lat, lon, settlement_data.get('elements', []))
# Query forests
forest_query = f"""
@@ -331,7 +383,19 @@ async def get_zone_features(lat: float, lon: float, direction: str, radius_m: in
);
out geom;
"""
forest_data = await query_overpass(forest_query)
# The four queries are independent - issue them concurrently.
roads_data, water_data, settlement_data, forest_data = await asyncio.gather(
query_overpass(roads_query),
query_overpass(water_query),
query_overpass(settlement_query),
query_overpass(forest_query),
)
roads_km = calculate_road_length(roads_data.get('elements', []))
water_distance = find_nearest_distance(lat, lon, water_data.get('elements', []))
settlement_distance = find_nearest_distance(lat, lon, settlement_data.get('elements', []))
forest_pct = calculate_forest_coverage(forest_data.get('elements', []), radius_m)
# Calculate road density (km of roads per km²)
+80
View File
@@ -0,0 +1,80 @@
"""
Statistical search-priority recommendation.
Single home for the recommendation scoring rules (B3). Previously duplicated
between `backend/routers/stats.py` (inline, exact-match on terrain/weather,
plus a health_flags rule) and `backend/services/stats_service.py` (substring
match on tolerant key aliases, no health_flags rule).
The unified rules below are the union of the two: identical weights and
threshold, substring matching (a superset of the old exact match), tolerant
input keys, and the health_flags rule kept. This is NOT the 7-factor zonal
scoring of `scoring_service` — it only labels a case high/normal priority.
Deliberately free of any database import so both the router and the service
layer can use it without side effects.
"""
from __future__ import annotations
from typing import Any
# Scoring weights and threshold — unchanged from both previous implementations.
WEIGHT_YOUNG_CHILD = 20
WEIGHT_LONG_ELAPSED = 20
WEIGHT_RISKY_TERRAIN = 15
WEIGHT_ADVERSE_WEATHER = 15
WEIGHT_MULTIPLE_HEALTH_FLAGS = 15
HIGH_PRIORITY_THRESHOLD = 40
YOUNG_CHILD_AGE = 12
LONG_ELAPSED_HOURS = 12
MULTIPLE_HEALTH_FLAGS = 2
RISKY_TERRAIN_TOKENS = ('лес', 'болото', 'вода')
ADVERSE_WEATHER_TOKENS = ('дождь', 'туман', 'снег', 'ночь')
HIGH_PRIORITY_TEXT = 'Высокий приоритет на прочёс и дрон'
NORMAL_PRIORITY_TEXT = 'Стандартный приоритет поиска'
def score_recommendation(
age: int | None = None,
elapsed_hours: int | None = None,
terrain: str | None = None,
weather: str | None = None,
health_flags: list[str] | None = None,
) -> dict[str, Any]:
"""Score a case and return {'score', 'priority', 'recommendation'}."""
score = 0
if age is not None and age < YOUNG_CHILD_AGE:
score += WEIGHT_YOUNG_CHILD
if elapsed_hours is not None and elapsed_hours >= LONG_ELAPSED_HOURS:
score += WEIGHT_LONG_ELAPSED
if any(token in str(terrain or '').lower() for token in RISKY_TERRAIN_TOKENS):
score += WEIGHT_RISKY_TERRAIN
if any(token in str(weather or '').lower() for token in ADVERSE_WEATHER_TOKENS):
score += WEIGHT_ADVERSE_WEATHER
if len(health_flags or []) >= MULTIPLE_HEALTH_FLAGS:
score += WEIGHT_MULTIPLE_HEALTH_FLAGS
is_high = score >= HIGH_PRIORITY_THRESHOLD
return {
'score': score,
'priority': 'high' if is_high else 'normal',
'recommendation': HIGH_PRIORITY_TEXT if is_high else NORMAL_PRIORITY_TEXT,
}
def get_statistical_recommendation(case_data: dict[str, Any]) -> dict[str, Any]:
"""Dict-based entry point, tolerant of the field aliases used across
the desktop / mobile / admin payloads."""
return score_recommendation(
age=case_data.get('age') if case_data.get('age') is not None else case_data.get('age_years'),
elapsed_hours=case_data.get('elapsed_hours'),
terrain=case_data.get('terrain') or case_data.get('terrain_primary'),
weather=case_data.get('weather') or case_data.get('precipitation'),
health_flags=case_data.get('health_flags'),
)