980 lines
36 KiB
Python
980 lines
36 KiB
Python
from __future__ import annotations
|
|
|
|
import json
|
|
import os
|
|
import threading
|
|
import time
|
|
import hashlib
|
|
from concurrent.futures import TimeoutError as FutureTimeoutError
|
|
from concurrent.futures import ThreadPoolExecutor, as_completed
|
|
from datetime import datetime, timedelta
|
|
from typing import Any, Dict, List, Optional
|
|
|
|
import httpx
|
|
from loguru import logger
|
|
|
|
from web.analysis_service import _analyze, _build_city_market_scan_payload
|
|
from web.core import CITIES
|
|
|
|
_SCAN_TERMINAL_CACHE_LOCK = threading.Lock()
|
|
_SCAN_TERMINAL_CACHE: Dict[str, Dict[str, Any]] = {}
|
|
_SCAN_TERMINAL_REFRESHING: set[str] = set()
|
|
_SCAN_TERMINAL_AI_CACHE_LOCK = threading.Lock()
|
|
_SCAN_TERMINAL_AI_CACHE: Dict[str, Dict[str, Any]] = {}
|
|
SCAN_TERMINAL_PAYLOAD_TTL_SEC = max(
|
|
5,
|
|
int(os.getenv("POLYWEATHER_SCAN_TERMINAL_PAYLOAD_TTL_SEC", "30")),
|
|
)
|
|
SCAN_TERMINAL_BUILD_TIMEOUT_SEC = max(
|
|
8,
|
|
int(os.getenv("POLYWEATHER_SCAN_TERMINAL_BUILD_TIMEOUT_SEC", "22")),
|
|
)
|
|
SCAN_AI_MODEL = str(
|
|
os.getenv("POLYWEATHER_SCAN_AI_MODEL") or "deepseek-v4-flash"
|
|
).strip()
|
|
SCAN_AI_BASE_URL = str(
|
|
os.getenv("POLYWEATHER_DEEPSEEK_BASE_URL") or "https://api.deepseek.com"
|
|
).strip().rstrip("/")
|
|
SCAN_AI_ENABLED = str(
|
|
os.getenv("POLYWEATHER_SCAN_AI_ENABLED") or "false"
|
|
).strip().lower() in {"1", "true", "yes", "on"}
|
|
SCAN_AI_TIMEOUT_SEC = max(
|
|
3,
|
|
int(os.getenv("POLYWEATHER_SCAN_AI_TIMEOUT_SEC", "12")),
|
|
)
|
|
SCAN_AI_CACHE_TTL_SEC = max(
|
|
30,
|
|
int(os.getenv("POLYWEATHER_SCAN_AI_CACHE_TTL_SEC", "600")),
|
|
)
|
|
SCAN_AI_MAX_ROWS = max(
|
|
1,
|
|
int(os.getenv("POLYWEATHER_SCAN_AI_MAX_ROWS", "40")),
|
|
)
|
|
|
|
|
|
def _safe_float(value: Any) -> Optional[float]:
|
|
try:
|
|
if value is None or value == "":
|
|
return None
|
|
return float(value)
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
def _safe_int(value: Any, default: int) -> int:
|
|
try:
|
|
return int(value)
|
|
except Exception:
|
|
return int(default)
|
|
|
|
|
|
def _normalize_scan_terminal_filters(
|
|
raw_filters: Optional[Dict[str, Any]] = None,
|
|
) -> Dict[str, Any]:
|
|
raw = raw_filters if isinstance(raw_filters, dict) else {}
|
|
min_price = _safe_float(raw.get("min_price"))
|
|
max_price = _safe_float(raw.get("max_price"))
|
|
if min_price is None:
|
|
min_price = 0.05
|
|
if max_price is None:
|
|
max_price = 0.95
|
|
min_price = max(0.0, min(1.0, min_price))
|
|
max_price = max(0.0, min(1.0, max_price))
|
|
if min_price > max_price:
|
|
min_price, max_price = max_price, min_price
|
|
|
|
high_liquidity_only = bool(raw.get("high_liquidity_only"))
|
|
min_liquidity = _safe_float(raw.get("min_liquidity"))
|
|
if min_liquidity is None:
|
|
min_liquidity = 5000.0 if high_liquidity_only else 500.0
|
|
if high_liquidity_only:
|
|
min_liquidity = max(min_liquidity, 5000.0)
|
|
|
|
return {
|
|
"scan_mode": str(raw.get("scan_mode") or "tradable").strip().lower()
|
|
or "tradable",
|
|
"min_price": float(min_price),
|
|
"max_price": float(max_price),
|
|
"min_edge_pct": max(0.0, _safe_float(raw.get("min_edge_pct")) or 2.0),
|
|
"min_liquidity": max(0.0, float(min_liquidity)),
|
|
"high_liquidity_only": high_liquidity_only,
|
|
"market_type": str(raw.get("market_type") or "maxtemp").strip().lower()
|
|
or "maxtemp",
|
|
"time_range": str(raw.get("time_range") or "today").strip().lower()
|
|
or "today",
|
|
"limit": max(1, min(_safe_int(raw.get("limit"), 25), 100)),
|
|
"max_spread": max(0.0, _safe_float(raw.get("max_spread")) or 0.03),
|
|
}
|
|
|
|
|
|
def _scan_terminal_cache_key(filters: Dict[str, Any]) -> str:
|
|
normalized = _normalize_scan_terminal_filters(filters)
|
|
return json.dumps(normalized, ensure_ascii=True, sort_keys=True)
|
|
|
|
|
|
def _get_cached_scan_terminal_payload(
|
|
filters: Dict[str, Any],
|
|
) -> Optional[Dict[str, Any]]:
|
|
cache_key = _scan_terminal_cache_key(filters)
|
|
now = time.time()
|
|
with _SCAN_TERMINAL_CACHE_LOCK:
|
|
cached = _SCAN_TERMINAL_CACHE.get(cache_key)
|
|
if not cached:
|
|
return None
|
|
cached_at = float(cached.get("t") or 0.0)
|
|
if now - cached_at >= float(SCAN_TERMINAL_PAYLOAD_TTL_SEC):
|
|
return None
|
|
payload = cached.get("payload")
|
|
if not isinstance(payload, dict):
|
|
return None
|
|
return dict(payload)
|
|
|
|
|
|
def _get_scan_terminal_cache_entry(filters: Dict[str, Any]) -> Optional[Dict[str, Any]]:
|
|
cache_key = _scan_terminal_cache_key(filters)
|
|
with _SCAN_TERMINAL_CACHE_LOCK:
|
|
cached = _SCAN_TERMINAL_CACHE.get(cache_key)
|
|
if not isinstance(cached, dict):
|
|
return None
|
|
return dict(cached)
|
|
|
|
|
|
def _build_scan_terminal_snapshot_id(
|
|
filters: Dict[str, Any],
|
|
rows: List[Dict[str, Any]],
|
|
summary: Dict[str, Any],
|
|
top_signal: Optional[Dict[str, Any]],
|
|
) -> str:
|
|
seed_payload = {
|
|
"filters": filters,
|
|
"summary": {
|
|
"candidate_total": summary.get("candidate_total"),
|
|
"tradable_market_count": summary.get("tradable_market_count"),
|
|
"avg_edge_percent": summary.get("avg_edge_percent"),
|
|
},
|
|
"top_signal": {
|
|
"id": (top_signal or {}).get("id"),
|
|
"edge_percent": (top_signal or {}).get("edge_percent"),
|
|
"final_score": (top_signal or {}).get("final_score"),
|
|
},
|
|
"rows": [
|
|
{
|
|
"id": row.get("id"),
|
|
"edge_percent": row.get("edge_percent"),
|
|
"final_score": row.get("final_score"),
|
|
}
|
|
for row in rows[:10]
|
|
],
|
|
}
|
|
digest = hashlib.md5(
|
|
json.dumps(seed_payload, ensure_ascii=True, sort_keys=True).encode("utf-8")
|
|
).hexdigest()
|
|
return f"scan-{digest[:10]}"
|
|
|
|
|
|
def _set_cached_scan_terminal_payload(
|
|
filters: Dict[str, Any],
|
|
payload: Dict[str, Any],
|
|
) -> None:
|
|
cache_key = _scan_terminal_cache_key(filters)
|
|
existing = _get_scan_terminal_cache_entry(filters) or {}
|
|
with _SCAN_TERMINAL_CACHE_LOCK:
|
|
_SCAN_TERMINAL_CACHE[cache_key] = {
|
|
"t": time.time(),
|
|
"payload": dict(payload),
|
|
"success_t": time.time(),
|
|
"success_payload": dict(payload),
|
|
"last_error": existing.get("last_error"),
|
|
"last_failed_at": existing.get("last_failed_at"),
|
|
}
|
|
|
|
|
|
def _set_scan_terminal_failure_state(
|
|
filters: Dict[str, Any],
|
|
*,
|
|
error_message: str,
|
|
) -> None:
|
|
cache_key = _scan_terminal_cache_key(filters)
|
|
with _SCAN_TERMINAL_CACHE_LOCK:
|
|
existing = _SCAN_TERMINAL_CACHE.get(cache_key) or {}
|
|
existing["last_error"] = error_message
|
|
existing["last_failed_at"] = datetime.utcnow().isoformat() + "Z"
|
|
_SCAN_TERMINAL_CACHE[cache_key] = existing
|
|
|
|
|
|
def _start_scan_terminal_background_refresh(filters: Dict[str, Any]) -> bool:
|
|
cache_key = _scan_terminal_cache_key(filters)
|
|
with _SCAN_TERMINAL_CACHE_LOCK:
|
|
if cache_key in _SCAN_TERMINAL_REFRESHING:
|
|
return False
|
|
_SCAN_TERMINAL_REFRESHING.add(cache_key)
|
|
|
|
def _runner() -> None:
|
|
try:
|
|
_build_scan_terminal_payload_uncached(filters, force_refresh=True)
|
|
except Exception as exc: # pragma: no cover - defensive background guard
|
|
logger.warning("scan terminal background refresh failed: {}", exc)
|
|
finally:
|
|
with _SCAN_TERMINAL_CACHE_LOCK:
|
|
_SCAN_TERMINAL_REFRESHING.discard(cache_key)
|
|
|
|
thread = threading.Thread(
|
|
target=_runner,
|
|
name="polyweather-scan-terminal-refresh",
|
|
daemon=True,
|
|
)
|
|
thread.start()
|
|
return True
|
|
|
|
|
|
def _build_stale_scan_terminal_payload(
|
|
*,
|
|
filters: Dict[str, Any],
|
|
success_payload: Dict[str, Any],
|
|
error_message: str,
|
|
failed_at: Optional[str],
|
|
) -> Dict[str, Any]:
|
|
payload = dict(success_payload)
|
|
payload["status"] = "stale"
|
|
payload["stale"] = True
|
|
payload["stale_reason"] = error_message
|
|
payload["last_success_at"] = success_payload.get("generated_at")
|
|
payload["last_failed_at"] = failed_at
|
|
payload["filters"] = filters
|
|
return payload
|
|
|
|
|
|
def _build_failed_scan_terminal_payload(
|
|
*,
|
|
filters: Dict[str, Any],
|
|
error_message: str,
|
|
failed_at: Optional[str] = None,
|
|
) -> Dict[str, Any]:
|
|
return {
|
|
"generated_at": datetime.utcnow().isoformat() + "Z",
|
|
"snapshot_id": None,
|
|
"status": "failed",
|
|
"stale": False,
|
|
"stale_reason": error_message,
|
|
"last_success_at": None,
|
|
"last_failed_at": failed_at or (datetime.utcnow().isoformat() + "Z"),
|
|
"filters": filters,
|
|
"summary": {
|
|
"recommended_count": 0,
|
|
"visible_count": 0,
|
|
"candidate_total": 0,
|
|
"avg_edge_percent": None,
|
|
"avg_primary_confidence": None,
|
|
"tradable_market_count": 0,
|
|
"total_volume": 0.0,
|
|
"resolved_market_type": "maxtemp",
|
|
},
|
|
"top_signal": None,
|
|
"rows": [],
|
|
}
|
|
|
|
|
|
def _extract_ai_json_object(raw_text: str) -> Dict[str, Any]:
|
|
text = str(raw_text or "").strip()
|
|
if not text:
|
|
raise ValueError("empty AI content")
|
|
try:
|
|
parsed = json.loads(text)
|
|
if isinstance(parsed, dict):
|
|
return parsed
|
|
except Exception:
|
|
pass
|
|
start = text.find("{")
|
|
end = text.rfind("}")
|
|
if start >= 0 and end > start:
|
|
parsed = json.loads(text[start : end + 1])
|
|
if isinstance(parsed, dict):
|
|
return parsed
|
|
raise ValueError("AI content is not a JSON object")
|
|
|
|
|
|
def _scan_ai_cache_key(snapshot_id: str, filters: Dict[str, Any]) -> str:
|
|
raw = json.dumps(
|
|
{
|
|
"snapshot_id": snapshot_id,
|
|
"filters": _normalize_scan_terminal_filters(filters),
|
|
"model": SCAN_AI_MODEL,
|
|
"max_rows": SCAN_AI_MAX_ROWS,
|
|
},
|
|
sort_keys=True,
|
|
ensure_ascii=False,
|
|
)
|
|
return hashlib.sha256(raw.encode("utf-8")).hexdigest()
|
|
|
|
|
|
def _get_cached_scan_ai_result(snapshot_id: str, filters: Dict[str, Any]) -> Optional[Dict[str, Any]]:
|
|
cache_key = _scan_ai_cache_key(snapshot_id, filters)
|
|
now = time.time()
|
|
with _SCAN_TERMINAL_AI_CACHE_LOCK:
|
|
cached = _SCAN_TERMINAL_AI_CACHE.get(cache_key)
|
|
if not cached:
|
|
return None
|
|
cached_at = float(cached.get("cached_at") or 0.0)
|
|
if now - cached_at >= float(SCAN_AI_CACHE_TTL_SEC):
|
|
return None
|
|
result = cached.get("result")
|
|
if isinstance(result, dict):
|
|
return dict(result)
|
|
return None
|
|
|
|
|
|
def _set_cached_scan_ai_result(snapshot_id: str, filters: Dict[str, Any], result: Dict[str, Any]) -> None:
|
|
cache_key = _scan_ai_cache_key(snapshot_id, filters)
|
|
with _SCAN_TERMINAL_AI_CACHE_LOCK:
|
|
_SCAN_TERMINAL_AI_CACHE[cache_key] = {
|
|
"cached_at": time.time(),
|
|
"result": result,
|
|
}
|
|
|
|
|
|
def _compact_ai_candidate(row: Dict[str, Any]) -> Dict[str, Any]:
|
|
return {
|
|
"id": row.get("id"),
|
|
"city": row.get("city_display_name") or row.get("display_name") or row.get("city"),
|
|
"local_time": row.get("local_time"),
|
|
"current_temp": row.get("current_temp"),
|
|
"current_max_so_far": row.get("current_max_so_far"),
|
|
"deb_prediction": row.get("deb_prediction"),
|
|
"action": row.get("action"),
|
|
"side": row.get("side"),
|
|
"target_label": row.get("target_label"),
|
|
"target_value": row.get("target_value"),
|
|
"target_threshold": row.get("target_threshold"),
|
|
"target_unit": row.get("target_unit"),
|
|
"model_probability": row.get("model_probability"),
|
|
"market_probability": row.get("market_probability"),
|
|
"model_event_probability": row.get("model_event_probability"),
|
|
"market_event_probability": row.get("market_event_probability"),
|
|
"yes_ask": row.get("yes_ask"),
|
|
"no_ask": row.get("no_ask"),
|
|
"ask": row.get("ask"),
|
|
"spread": row.get("spread"),
|
|
"quote_age_ms": row.get("quote_age_ms"),
|
|
"edge_percent": row.get("edge_percent"),
|
|
"final_score": row.get("final_score"),
|
|
"window_phase": row.get("window_phase"),
|
|
"remaining_window_minutes": row.get("remaining_window_minutes"),
|
|
"distribution_preview": (row.get("distribution_preview") or [])[:6],
|
|
"peak_value": row.get("peak_value"),
|
|
"peak_probability": row.get("peak_probability"),
|
|
"cluster_role": row.get("cluster_role"),
|
|
"cluster_core_low": row.get("cluster_core_low"),
|
|
"cluster_core_high": row.get("cluster_core_high"),
|
|
"cluster_median": row.get("cluster_median"),
|
|
"cluster_deb_reference": row.get("cluster_deb_reference"),
|
|
"cluster_model_count": row.get("cluster_model_count"),
|
|
"trend_alignment": row.get("trend_alignment"),
|
|
"tradable": row.get("tradable"),
|
|
"accepting_orders": row.get("accepting_orders"),
|
|
}
|
|
|
|
|
|
def _build_scan_ai_prompt(payload: Dict[str, Any]) -> Dict[str, Any]:
|
|
rows = [
|
|
_compact_ai_candidate(row)
|
|
for row in (payload.get("rows") or [])[:SCAN_AI_MAX_ROWS]
|
|
if isinstance(row, dict) and row.get("id")
|
|
]
|
|
return {
|
|
"snapshot_id": payload.get("snapshot_id"),
|
|
"generated_at": payload.get("generated_at"),
|
|
"summary": payload.get("summary") or {},
|
|
"filters": payload.get("filters") or {},
|
|
"candidates": rows,
|
|
}
|
|
|
|
|
|
def _call_deepseek_scan_ai(ai_input: Dict[str, Any]) -> Dict[str, Any]:
|
|
api_key = str(os.getenv("POLYWEATHER_DEEPSEEK_API_KEY") or "").strip()
|
|
if not api_key:
|
|
raise RuntimeError("POLYWEATHER_DEEPSEEK_API_KEY is not configured")
|
|
|
|
system_prompt = (
|
|
"你是 PolyWeather 的付费市场扫描员。你只能基于用户提供的 JSON 快照做判断,"
|
|
"不得编造城市、价格、概率、盘口或天气数据。你的任务是从候选 rows 中选择真正值得看的机会,"
|
|
"尤其要检查 DEB、EMOS 分布、模型集群、当前实测、峰值窗口和盘口价格是否一致。"
|
|
"若模型集群不支持某个 YES,不要机械推荐 YES;可以 veto 或改为偏向 NO 候选。"
|
|
"只能引用输入中的 row id。必须输出 JSON object,字段为 recommendations、vetoed、downgraded、summary_zh、summary_en。"
|
|
)
|
|
user_payload = {
|
|
"task": (
|
|
"Review these existing PolyWeather market candidates. "
|
|
"Return strict JSON only. recommendations items need row_id, rank, decision, confidence, "
|
|
"reason_zh, reason_en, model_cluster_note. vetoed/downgraded items need row_id and reason."
|
|
),
|
|
"snapshot": ai_input,
|
|
}
|
|
with httpx.Client(timeout=float(SCAN_AI_TIMEOUT_SEC)) as client:
|
|
response = client.post(
|
|
f"{SCAN_AI_BASE_URL}/chat/completions",
|
|
headers={
|
|
"Authorization": f"Bearer {api_key}",
|
|
"Content-Type": "application/json",
|
|
},
|
|
json={
|
|
"model": SCAN_AI_MODEL,
|
|
"temperature": 0.1,
|
|
"max_tokens": 2200,
|
|
"response_format": {"type": "json_object"},
|
|
"messages": [
|
|
{"role": "system", "content": system_prompt},
|
|
{
|
|
"role": "user",
|
|
"content": json.dumps(user_payload, ensure_ascii=False),
|
|
},
|
|
],
|
|
},
|
|
)
|
|
response.raise_for_status()
|
|
data = response.json()
|
|
content = (
|
|
((data.get("choices") or [{}])[0].get("message") or {}).get("content")
|
|
if isinstance(data, dict)
|
|
else None
|
|
)
|
|
return _extract_ai_json_object(str(content or ""))
|
|
|
|
|
|
def _normalize_ai_items(raw_items: Any) -> List[Dict[str, Any]]:
|
|
if not isinstance(raw_items, list):
|
|
return []
|
|
out: List[Dict[str, Any]] = []
|
|
for item in raw_items:
|
|
if isinstance(item, str):
|
|
out.append({"row_id": item})
|
|
elif isinstance(item, dict):
|
|
row_id = str(item.get("row_id") or item.get("id") or "").strip()
|
|
if row_id:
|
|
out.append({**item, "row_id": row_id})
|
|
return out
|
|
|
|
|
|
def _merge_scan_ai_result(payload: Dict[str, Any], ai_raw: Dict[str, Any], *, cached: bool = False) -> Dict[str, Any]:
|
|
rows = [dict(row) for row in (payload.get("rows") or []) if isinstance(row, dict)]
|
|
by_id = {str(row.get("id")): row for row in rows if row.get("id")}
|
|
recommendations = _normalize_ai_items(ai_raw.get("recommendations"))
|
|
vetoed = _normalize_ai_items(ai_raw.get("vetoed"))
|
|
downgraded = _normalize_ai_items(ai_raw.get("downgraded"))
|
|
|
|
veto_ids = {str(item.get("row_id")) for item in vetoed}
|
|
downgrade_ids = {str(item.get("row_id")) for item in downgraded}
|
|
recommended_ids: set[str] = set()
|
|
|
|
for item in vetoed:
|
|
row = by_id.get(str(item.get("row_id")))
|
|
if not row:
|
|
continue
|
|
row["ai_decision"] = "veto"
|
|
row["ai_reason_zh"] = item.get("reason_zh") or item.get("reason")
|
|
row["ai_reason_en"] = item.get("reason_en")
|
|
for item in downgraded:
|
|
row = by_id.get(str(item.get("row_id")))
|
|
if not row:
|
|
continue
|
|
row["ai_decision"] = "downgrade"
|
|
row["ai_reason_zh"] = item.get("reason_zh") or item.get("reason")
|
|
row["ai_reason_en"] = item.get("reason_en")
|
|
for fallback_rank, item in enumerate(recommendations, start=1):
|
|
row_id = str(item.get("row_id"))
|
|
row = by_id.get(row_id)
|
|
if not row:
|
|
continue
|
|
if row_id in veto_ids:
|
|
continue
|
|
recommended_ids.add(row_id)
|
|
row["ai_decision"] = str(item.get("decision") or "approve").strip().lower() or "approve"
|
|
row["ai_rank"] = _safe_int(item.get("rank"), fallback_rank)
|
|
row["ai_confidence"] = item.get("confidence")
|
|
row["ai_reason_zh"] = item.get("reason_zh") or item.get("reason")
|
|
row["ai_reason_en"] = item.get("reason_en")
|
|
row["ai_model_cluster_note"] = item.get("model_cluster_note")
|
|
|
|
for row in rows:
|
|
row_id = str(row.get("id"))
|
|
if row_id not in recommended_ids and row_id not in veto_ids and row_id not in downgrade_ids:
|
|
row["ai_decision"] = row.get("ai_decision") or "neutral"
|
|
|
|
def _ai_sort_key(row: Dict[str, Any]) -> tuple:
|
|
decision = str(row.get("ai_decision") or "").lower()
|
|
if decision == "veto":
|
|
tier = 3
|
|
elif decision == "downgrade":
|
|
tier = 2
|
|
elif row.get("ai_rank") is not None:
|
|
tier = 0
|
|
else:
|
|
tier = 1
|
|
return (
|
|
tier,
|
|
_safe_int(row.get("ai_rank"), 999),
|
|
-float(row.get("final_score") or 0.0),
|
|
-float(row.get("edge_percent") or 0.0),
|
|
)
|
|
|
|
rows.sort(key=_ai_sort_key)
|
|
top_signal = next(
|
|
(row for row in rows if str(row.get("ai_decision") or "").lower() != "veto"),
|
|
rows[0] if rows else None,
|
|
)
|
|
ai_scan = {
|
|
"status": "ready",
|
|
"model": SCAN_AI_MODEL,
|
|
"cached": cached,
|
|
"generated_at": datetime.utcnow().isoformat() + "Z",
|
|
"summary_zh": ai_raw.get("summary_zh"),
|
|
"summary_en": ai_raw.get("summary_en"),
|
|
"recommended_count": sum(1 for row in rows if row.get("ai_rank") is not None),
|
|
"vetoed_count": sum(1 for row in rows if row.get("ai_decision") == "veto"),
|
|
"downgraded_count": sum(1 for row in rows if row.get("ai_decision") == "downgrade"),
|
|
}
|
|
merged = {
|
|
**payload,
|
|
"rows": rows,
|
|
"top_signal": top_signal,
|
|
"ai_scan": ai_scan,
|
|
}
|
|
return merged
|
|
|
|
|
|
def _build_scan_ai_unavailable_payload(
|
|
payload: Dict[str, Any],
|
|
*,
|
|
status: str,
|
|
reason: str,
|
|
) -> Dict[str, Any]:
|
|
return {
|
|
**payload,
|
|
"ai_scan": {
|
|
"status": status,
|
|
"model": SCAN_AI_MODEL,
|
|
"cached": False,
|
|
"generated_at": datetime.utcnow().isoformat() + "Z",
|
|
"reason": reason,
|
|
},
|
|
}
|
|
|
|
|
|
def _resolve_time_range_dates(data: Dict[str, Any], time_range: str) -> List[str]:
|
|
local_date = str(data.get("local_date") or "").strip()
|
|
multi_model_daily = data.get("multi_model_daily") or {}
|
|
available_dates = sorted(
|
|
str(date_key).strip()
|
|
for date_key in (multi_model_daily.keys() if isinstance(multi_model_daily, dict) else [])
|
|
if str(date_key).strip()
|
|
)
|
|
|
|
if not local_date:
|
|
return available_dates[:1]
|
|
if time_range == "today":
|
|
return [local_date]
|
|
|
|
try:
|
|
local_dt = datetime.fromisoformat(local_date)
|
|
except Exception:
|
|
return available_dates[:7] if time_range == "week" else available_dates[:1]
|
|
|
|
if time_range == "tomorrow":
|
|
target = (local_dt + timedelta(days=1)).strftime("%Y-%m-%d")
|
|
if target in available_dates:
|
|
return [target]
|
|
future_dates = [date_key for date_key in available_dates if date_key > local_date]
|
|
return future_dates[:1]
|
|
|
|
if time_range == "week":
|
|
target_dates = [date_key for date_key in available_dates if date_key >= local_date]
|
|
if local_date not in target_dates:
|
|
target_dates.insert(0, local_date)
|
|
deduped: List[str] = []
|
|
for date_key in target_dates:
|
|
if date_key not in deduped:
|
|
deduped.append(date_key)
|
|
if len(deduped) >= 7:
|
|
break
|
|
return deduped
|
|
|
|
return [local_date]
|
|
|
|
|
|
def _build_terminal_row(
|
|
*,
|
|
city: str,
|
|
data: Dict[str, Any],
|
|
scan: Dict[str, Any],
|
|
row: Dict[str, Any],
|
|
) -> Dict[str, Any]:
|
|
current = data.get("current") or {}
|
|
multi_model_daily = data.get("multi_model_daily") or {}
|
|
selected_date = str(row.get("selected_date") or scan.get("selected_date") or data.get("local_date") or "").strip()
|
|
daily_entry = multi_model_daily.get(selected_date) if isinstance(multi_model_daily, dict) else {}
|
|
if not isinstance(daily_entry, dict):
|
|
daily_entry = {}
|
|
|
|
display_name = str(data.get("display_name") or city).strip() or city
|
|
market_slug = str(row.get("market_slug") or "").strip()
|
|
side = str(row.get("side") or "").strip().lower()
|
|
edge_percent = _safe_float(row.get("edge_percent"))
|
|
final_score = _safe_float(row.get("final_score"))
|
|
volume = _safe_float(row.get("volume")) or 0.0
|
|
primary_signal = scan.get("primary_signal") or {}
|
|
|
|
return {
|
|
**row,
|
|
"id": str(row.get("id") or f"{city}|{selected_date}|{market_slug}|{side}"),
|
|
"city": city,
|
|
"city_display_name": display_name,
|
|
"selected_date": selected_date or None,
|
|
"local_date": data.get("local_date"),
|
|
"local_time": data.get("local_time"),
|
|
"temp_symbol": data.get("temp_symbol"),
|
|
"current_temp": current.get("temp"),
|
|
"current_max_so_far": current.get("max_so_far"),
|
|
"deb_prediction": ((daily_entry.get("deb") or {}).get("prediction") if isinstance(daily_entry.get("deb"), dict) else None)
|
|
or ((data.get("deb") or {}).get("prediction") if isinstance(data.get("deb"), dict) else None),
|
|
"display_name": display_name,
|
|
"airport": ((data.get("risk") or {}).get("airport") if isinstance(data.get("risk"), dict) else None),
|
|
"risk_level": ((data.get("risk") or {}).get("level") if isinstance(data.get("risk"), dict) else None),
|
|
"distribution_bias": scan.get("distribution_bias"),
|
|
"distribution_preview": scan.get("distribution_preview") or row.get("distribution_preview") or [],
|
|
"window_phase": row.get("window_phase") or scan.get("window_phase"),
|
|
"window_score": row.get("window_score") if row.get("window_score") is not None else scan.get("window_score"),
|
|
"signal_status": scan.get("signal_status"),
|
|
"candidate_count": scan.get("candidate_count"),
|
|
"resolved_market_type": scan.get("resolved_market_type") or "maxtemp",
|
|
"market_key": f"{city}|{selected_date}|{market_slug}",
|
|
"is_primary_signal": bool(primary_signal and primary_signal.get("id") == row.get("id")),
|
|
"signal_confidence": final_score,
|
|
"edge_percent": edge_percent,
|
|
"final_score": final_score,
|
|
"volume": volume,
|
|
}
|
|
|
|
|
|
def _scan_city_terminal_rows(
|
|
city: str,
|
|
filters: Dict[str, Any],
|
|
*,
|
|
force_refresh: bool = False,
|
|
) -> Dict[str, Any]:
|
|
data = _analyze(
|
|
city,
|
|
force_refresh=force_refresh,
|
|
include_llm_commentary=False,
|
|
detail_mode="market",
|
|
)
|
|
target_dates = _resolve_time_range_dates(data, filters["time_range"])
|
|
rows: List[Dict[str, Any]] = []
|
|
primary_scores: List[float] = []
|
|
candidate_total = 0
|
|
|
|
for target_date in target_dates:
|
|
payload = _build_city_market_scan_payload(
|
|
data,
|
|
market_slug=None,
|
|
target_date=target_date,
|
|
lite=True,
|
|
scan_filters=filters,
|
|
)
|
|
scan = payload.get("market_scan") or {}
|
|
candidate_total += int(scan.get("candidate_count") or 0)
|
|
raw_rows = scan.get("scan_rows")
|
|
if not isinstance(raw_rows, list) or not raw_rows:
|
|
raw_rows = [scan.get("primary_signal")] if isinstance(scan.get("primary_signal"), dict) else []
|
|
if not raw_rows:
|
|
continue
|
|
for raw_row in raw_rows:
|
|
if not isinstance(raw_row, dict) or not raw_row:
|
|
continue
|
|
row = _build_terminal_row(
|
|
city=city,
|
|
data=data,
|
|
scan=scan,
|
|
row=raw_row,
|
|
)
|
|
rows.append(row)
|
|
score = _safe_float(row.get("final_score"))
|
|
if score is not None and row.get("is_primary_signal"):
|
|
primary_scores.append(score)
|
|
|
|
return {
|
|
"city": city,
|
|
"rows": rows,
|
|
"candidate_total": candidate_total,
|
|
"primary_scores": primary_scores,
|
|
}
|
|
|
|
|
|
def _build_scan_terminal_payload_uncached(
|
|
filters: Dict[str, Any],
|
|
*,
|
|
force_refresh: bool = False,
|
|
) -> Dict[str, Any]:
|
|
cached_entry = _get_scan_terminal_cache_entry(filters) or {}
|
|
|
|
try:
|
|
city_names = list(CITIES.keys())
|
|
max_workers = max(1, min(4, len(city_names)))
|
|
city_results: List[Dict[str, Any]] = []
|
|
failed_cities: List[str] = []
|
|
failed_reasons: List[str] = []
|
|
|
|
timed_out = False
|
|
timeout_message: Optional[str] = None
|
|
executor = ThreadPoolExecutor(max_workers=max_workers)
|
|
future_map = {
|
|
executor.submit(
|
|
_scan_city_terminal_rows,
|
|
city_name,
|
|
filters,
|
|
force_refresh=force_refresh,
|
|
): city_name
|
|
for city_name in city_names
|
|
}
|
|
try:
|
|
try:
|
|
completed = as_completed(
|
|
future_map,
|
|
timeout=float(SCAN_TERMINAL_BUILD_TIMEOUT_SEC),
|
|
)
|
|
for future in completed:
|
|
city_name = future_map[future]
|
|
try:
|
|
city_results.append(future.result())
|
|
except Exception as exc:
|
|
failed_cities.append(city_name)
|
|
failed_reasons.append(str(exc))
|
|
logger.warning("scan terminal city failed city={}: {}", city_name, exc)
|
|
except FutureTimeoutError:
|
|
timed_out = True
|
|
timeout_message = (
|
|
f"scan terminal build timed out after "
|
|
f"{SCAN_TERMINAL_BUILD_TIMEOUT_SEC}s"
|
|
)
|
|
failed_reasons.append(timeout_message)
|
|
for future, city_name in future_map.items():
|
|
if not future.done():
|
|
future.cancel()
|
|
failed_cities.append(city_name)
|
|
logger.warning(
|
|
"{}; completed={}/{}",
|
|
timeout_message,
|
|
len(city_results),
|
|
len(city_names),
|
|
)
|
|
finally:
|
|
executor.shutdown(wait=False)
|
|
|
|
if city_names and len(failed_cities) >= len(city_names):
|
|
error_message = failed_reasons[0] if failed_reasons else "all city market scans failed"
|
|
_set_scan_terminal_failure_state(filters, error_message=error_message)
|
|
failed_entry = _get_scan_terminal_cache_entry(filters) or {}
|
|
success_payload = failed_entry.get("success_payload")
|
|
failed_at = failed_entry.get("last_failed_at")
|
|
if isinstance(success_payload, dict) and success_payload:
|
|
return _build_stale_scan_terminal_payload(
|
|
filters=filters,
|
|
success_payload=success_payload,
|
|
error_message=error_message,
|
|
failed_at=failed_at,
|
|
)
|
|
return _build_failed_scan_terminal_payload(
|
|
filters=filters,
|
|
error_message=error_message,
|
|
failed_at=failed_at,
|
|
)
|
|
|
|
primary_rows: List[Dict[str, Any]] = []
|
|
primary_scores: List[float] = []
|
|
candidate_total = 0
|
|
|
|
for result in city_results:
|
|
candidate_total += int(result.get("candidate_total") or 0)
|
|
primary_rows.extend(result.get("rows") or [])
|
|
primary_scores.extend(result.get("primary_scores") or [])
|
|
|
|
primary_rows.sort(
|
|
key=lambda row: (
|
|
float(row.get("final_score") or 0.0),
|
|
float(row.get("edge_percent") or 0.0),
|
|
),
|
|
reverse=True,
|
|
)
|
|
|
|
ranked_rows: List[Dict[str, Any]] = []
|
|
for index, row in enumerate(primary_rows[: filters["limit"]], start=1):
|
|
ranked_rows.append(
|
|
{
|
|
**row,
|
|
"rank": index,
|
|
}
|
|
)
|
|
|
|
unique_market_volume: Dict[str, float] = {}
|
|
for row in primary_rows:
|
|
market_key = str(row.get("market_key") or row.get("id") or "").strip()
|
|
if not market_key:
|
|
continue
|
|
unique_market_volume[market_key] = max(
|
|
unique_market_volume.get(market_key, 0.0),
|
|
float(row.get("volume") or 0.0),
|
|
)
|
|
|
|
avg_edge = None
|
|
if primary_rows:
|
|
edge_values = [
|
|
float(row.get("edge_percent") or 0.0)
|
|
for row in primary_rows
|
|
if _safe_float(row.get("edge_percent")) is not None
|
|
]
|
|
if edge_values:
|
|
avg_edge = sum(edge_values) / len(edge_values)
|
|
|
|
avg_confidence = None
|
|
if primary_scores:
|
|
avg_confidence = sum(primary_scores) / len(primary_scores)
|
|
|
|
top_signal = ranked_rows[0] if ranked_rows else None
|
|
summary = {
|
|
"recommended_count": len(primary_rows),
|
|
"visible_count": len(ranked_rows),
|
|
"candidate_total": candidate_total,
|
|
"avg_edge_percent": avg_edge,
|
|
"avg_primary_confidence": avg_confidence,
|
|
"tradable_market_count": len(unique_market_volume),
|
|
"total_volume": sum(unique_market_volume.values()),
|
|
"resolved_market_type": "maxtemp",
|
|
"total_city_count": len(city_names),
|
|
"scanned_city_count": len(city_results),
|
|
"failed_city_count": len(failed_cities),
|
|
}
|
|
payload = {
|
|
"generated_at": datetime.utcnow().isoformat() + "Z",
|
|
"filters": filters,
|
|
"summary": summary,
|
|
"top_signal": top_signal,
|
|
"rows": ranked_rows,
|
|
"status": "partial" if timed_out else "ready",
|
|
"stale": False,
|
|
"stale_reason": timeout_message,
|
|
"last_success_at": None,
|
|
"last_failed_at": None,
|
|
}
|
|
payload["snapshot_id"] = _build_scan_terminal_snapshot_id(
|
|
filters,
|
|
ranked_rows,
|
|
summary,
|
|
top_signal,
|
|
)
|
|
|
|
_set_cached_scan_terminal_payload(filters, payload)
|
|
return payload
|
|
except Exception as exc:
|
|
error_message = str(exc)
|
|
logger.exception("scan terminal payload build failed: {}", error_message)
|
|
_set_scan_terminal_failure_state(filters, error_message=error_message)
|
|
success_payload = cached_entry.get("success_payload")
|
|
failed_at = _get_scan_terminal_cache_entry(filters).get("last_failed_at") if _get_scan_terminal_cache_entry(filters) else None
|
|
if isinstance(success_payload, dict) and success_payload:
|
|
return _build_stale_scan_terminal_payload(
|
|
filters=filters,
|
|
success_payload=success_payload,
|
|
error_message=error_message,
|
|
failed_at=failed_at,
|
|
)
|
|
return _build_failed_scan_terminal_payload(
|
|
filters=filters,
|
|
error_message=error_message,
|
|
failed_at=failed_at,
|
|
)
|
|
|
|
|
|
def build_scan_terminal_payload(
|
|
raw_filters: Optional[Dict[str, Any]] = None,
|
|
*,
|
|
force_refresh: bool = False,
|
|
) -> Dict[str, Any]:
|
|
filters = _normalize_scan_terminal_filters(raw_filters)
|
|
if not force_refresh:
|
|
cached = _get_cached_scan_terminal_payload(filters)
|
|
if cached is not None:
|
|
return cached
|
|
|
|
cached_entry = _get_scan_terminal_cache_entry(filters) or {}
|
|
success_payload = cached_entry.get("success_payload")
|
|
if isinstance(success_payload, dict) and success_payload:
|
|
started = _start_scan_terminal_background_refresh(filters)
|
|
return _build_stale_scan_terminal_payload(
|
|
filters=filters,
|
|
success_payload=success_payload,
|
|
error_message=(
|
|
"正在后台刷新市场扫描快照"
|
|
if started
|
|
else "市场扫描快照正在刷新中"
|
|
),
|
|
failed_at=cached_entry.get("last_failed_at"),
|
|
)
|
|
|
|
return _build_scan_terminal_payload_uncached(filters, force_refresh=force_refresh)
|
|
|
|
|
|
def build_scan_terminal_ai_payload(
|
|
raw_filters: Optional[Dict[str, Any]] = None,
|
|
*,
|
|
snapshot_id: Optional[str] = None,
|
|
) -> Dict[str, Any]:
|
|
filters = _normalize_scan_terminal_filters(raw_filters)
|
|
payload = build_scan_terminal_payload(filters, force_refresh=False)
|
|
current_snapshot_id = str(payload.get("snapshot_id") or "").strip()
|
|
requested_snapshot_id = str(snapshot_id or "").strip()
|
|
if requested_snapshot_id and current_snapshot_id and requested_snapshot_id != current_snapshot_id:
|
|
return _build_scan_ai_unavailable_payload(
|
|
payload,
|
|
status="snapshot_mismatch",
|
|
reason="scan snapshot changed; refresh the scan before running AI review",
|
|
)
|
|
if not current_snapshot_id:
|
|
return _build_scan_ai_unavailable_payload(
|
|
payload,
|
|
status="no_snapshot",
|
|
reason="no scan snapshot is available for AI review",
|
|
)
|
|
if not payload.get("rows"):
|
|
return _build_scan_ai_unavailable_payload(
|
|
payload,
|
|
status="no_rows",
|
|
reason="no candidate rows are available for AI review",
|
|
)
|
|
if not SCAN_AI_ENABLED:
|
|
return _build_scan_ai_unavailable_payload(
|
|
payload,
|
|
status="disabled",
|
|
reason="POLYWEATHER_SCAN_AI_ENABLED is not enabled",
|
|
)
|
|
if not str(os.getenv("POLYWEATHER_DEEPSEEK_API_KEY") or "").strip():
|
|
return _build_scan_ai_unavailable_payload(
|
|
payload,
|
|
status="missing_key",
|
|
reason="POLYWEATHER_DEEPSEEK_API_KEY is not configured",
|
|
)
|
|
|
|
cached = _get_cached_scan_ai_result(current_snapshot_id, filters)
|
|
if cached is not None:
|
|
return _merge_scan_ai_result(payload, cached, cached=True)
|
|
|
|
try:
|
|
ai_input = _build_scan_ai_prompt(payload)
|
|
ai_raw = _call_deepseek_scan_ai(ai_input)
|
|
_set_cached_scan_ai_result(current_snapshot_id, filters, ai_raw)
|
|
return _merge_scan_ai_result(payload, ai_raw, cached=False)
|
|
except Exception as exc:
|
|
logger.warning("scan terminal AI review failed: {}", exc)
|
|
return _build_scan_ai_unavailable_payload(
|
|
payload,
|
|
status="failed",
|
|
reason=str(exc),
|
|
)
|