Clarify scan terminal service boundaries

The scan terminal service had accumulated cache, payload, filtering, AI prompt, AI merge, METAR gate, ranking, and city-row construction details in one file. This splits those stable responsibilities into focused modules while preserving the endpoint payload shape and existing behavior.

Constraint: User-visible behavior and release version must remain unchanged for this internal refactor

Rejected: Rewrite the terminal scan flow around a new abstraction | too risky while production behavior is being stabilized

Confidence: high

Scope-risk: moderate

Directive: Keep scan_terminal_service.py as orchestration; add detailed rule changes to the focused modules instead of re-growing the service file

Tested: py_compile for extracted modules; ruff check .; pytest tests/test_scan_terminal_modules.py tests/test_web_observability.py; full pytest; npm run test:business; npm run build; git diff --cached --check
This commit is contained in:
2569718930@qq.com
2026-04-28 17:13:45 +08:00
parent b122e7cbae
commit 20516000a2
11 changed files with 1696 additions and 1309 deletions
+199
View File
@@ -0,0 +1,199 @@
from web.scan_terminal_ai_merge import merge_scan_ai_result
from web.scan_terminal_filters import normalize_scan_terminal_filters
from web.scan_terminal_metar_gate import _apply_metar_gate_to_row
from web.scan_terminal_payloads import (
build_failed_scan_terminal_payload,
build_scan_terminal_snapshot_id,
build_stale_scan_terminal_payload,
)
from web.scan_terminal_ranker import build_ranked_scan_terminal_result
def test_normalize_scan_terminal_filters_clamps_and_swaps_bounds():
filters = normalize_scan_terminal_filters(
{
"min_price": 1.2,
"max_price": -0.2,
"limit": 999,
"high_liquidity_only": True,
"min_liquidity": 100,
}
)
assert filters["min_price"] == 0.0
assert filters["max_price"] == 1.0
assert filters["limit"] == 100
assert filters["min_liquidity"] == 5000.0
def test_ranked_scan_terminal_result_sorts_and_summarizes_unique_markets():
result = build_ranked_scan_terminal_result(
city_results=[
{
"candidate_total": 2,
"primary_scores": [80.0],
"rows": [
{
"id": "low",
"market_key": "m1",
"final_score": 70.0,
"edge_percent": 4.0,
"volume": 100,
},
{
"id": "high",
"market_key": "m1",
"final_score": 90.0,
"edge_percent": 2.0,
"volume": 250,
},
],
},
{
"candidate_total": 1,
"primary_scores": [60.0],
"rows": [
{
"id": "tie-break",
"market_key": "m2",
"final_score": 90.0,
"edge_percent": 5.0,
"volume": 300,
}
],
},
],
filters={"limit": 2},
total_city_count=3,
failed_city_count=1,
)
assert [row["id"] for row in result["ranked_rows"]] == ["tie-break", "high"]
assert [row["rank"] for row in result["ranked_rows"]] == [1, 2]
assert result["top_signal"]["id"] == "tie-break"
assert result["summary"]["candidate_total"] == 3
assert result["summary"]["visible_count"] == 2
assert result["summary"]["tradable_market_count"] == 2
assert result["summary"]["total_volume"] == 550
assert result["summary"]["failed_city_count"] == 1
def test_scan_terminal_payload_helpers_preserve_stale_and_failed_shape():
success_payload = {
"generated_at": "2026-04-28T00:00:00Z",
"snapshot_id": "scan-old",
"filters": {"scan_mode": "tradable"},
"rows": [{"id": "row-1"}],
}
stale = build_stale_scan_terminal_payload(
filters={"scan_mode": "trend"},
success_payload=success_payload,
error_message="refresh failed",
failed_at="2026-04-28T00:01:00Z",
)
failed = build_failed_scan_terminal_payload(
filters={"scan_mode": "trend"},
error_message="network down",
failed_at="2026-04-28T00:02:00Z",
)
assert stale["status"] == "stale"
assert stale["stale"] is True
assert stale["rows"] == [{"id": "row-1"}]
assert stale["filters"] == {"scan_mode": "trend"}
assert stale["last_success_at"] == "2026-04-28T00:00:00Z"
assert failed["status"] == "failed"
assert failed["summary"]["candidate_total"] == 0
assert failed["rows"] == []
def test_scan_terminal_snapshot_id_is_stable_for_same_ranked_inputs():
summary = {
"candidate_total": 2,
"tradable_market_count": 2,
"avg_edge_percent": 3.5,
}
rows = [
{"id": "a", "edge_percent": 4.0, "final_score": 90.0},
{"id": "b", "edge_percent": 3.0, "final_score": 80.0},
]
first = build_scan_terminal_snapshot_id({"limit": 2}, rows, summary, rows[0])
second = build_scan_terminal_snapshot_id({"limit": 2}, rows, summary, rows[0])
assert first == second
assert first.startswith("scan-")
def test_metar_gate_vetoes_yes_when_observed_breaks_above_bucket():
row = {
"id": "yes-row",
"side": "yes",
"target_lower": 32.0,
"target_upper": 34.0,
"target_unit": "°C",
"metar_context": {
"obs_count": 6,
"max_temp": 35.0,
"last_temp": 35.0,
"trend_delta": 1.0,
"stale_for_today": False,
},
}
_apply_metar_gate_to_row(row)
assert row["v4_metar_decision"] == "veto"
assert row["ai_decision"] == "veto"
assert "越过目标桶上沿" in row["ai_reason_zh"]
def test_merge_scan_ai_result_applies_city_forecast_and_metadata():
payload = {
"snapshot_id": "scan-abc",
"rows": [
{
"id": "row-1",
"city": "manila",
"city_display_name": "Manila",
"final_score": 80.0,
"edge_percent": 5.0,
"metar_context": {"obs_count": 0},
}
],
}
ai_raw = {
"summary_zh": "摘要",
"city_forecasts": [
{
"city": "Manila",
"predicted_max": 35.0,
"range_low": 34.0,
"range_high": 36.0,
"reasoning_zh": "实测突破",
}
],
"recommendations": [{"row_id": "row-1", "rank": 1, "reason_zh": "观察"}],
"_polyweather_input_meta": {"sent_cities": 1, "sent_contracts": 1},
}
merged = merge_scan_ai_result(
payload,
ai_raw,
model="deepseek-v4-flash",
max_rows=40,
timeout_sec=40,
cache_ttl_sec=1800,
base_url="https://api.deepseek.com",
duration_ms=123,
input_rows=1,
)
row = merged["rows"][0]
assert row["ai_predicted_max"] == 35.0
assert row["ai_forecast_reason_zh"] == "实测突破"
assert row["ai_decision"] == "approve"
assert merged["ai_scan"]["sent_cities"] == 1
assert merged["ai_scan"]["sent_rows"] == 1
assert merged["ai_scan"]["duration_ms"] == 123
+6 -5
View File
@@ -3,8 +3,9 @@ from fastapi.testclient import TestClient
from web.app import app
import web.routes as routes
import web.scan_terminal_cache as scan_terminal_cache
import web.scan_terminal_service as scan_terminal_service
from web.scan_terminal_service import _scan_terminal_cache_key
from web.scan_terminal_cache import scan_terminal_cache_key
from src.database.db_manager import DBManager
from src.database.runtime_state import TruthRecordRepository, TrainingFeatureRecordRepository
@@ -737,7 +738,7 @@ def test_ops_truth_history_returns_filtered_rows(monkeypatch):
def test_scan_terminal_service_returns_stale_payload_after_failed_refresh(monkeypatch):
filters = {"scan_mode": "tradable", "limit": 5}
normalized_filters = scan_terminal_service._normalize_scan_terminal_filters(filters)
scan_terminal_service._SCAN_TERMINAL_CACHE.clear()
scan_terminal_cache._SCAN_TERMINAL_CACHE.clear()
monkeypatch.setattr(
scan_terminal_service,
@@ -778,7 +779,7 @@ def test_scan_terminal_service_returns_stale_payload_after_failed_refresh(monkey
def test_scan_terminal_service_returns_failed_without_success_snapshot(monkeypatch):
filters = {"scan_mode": "tradable", "limit": 5}
scan_terminal_service._SCAN_TERMINAL_CACHE.clear()
scan_terminal_cache._SCAN_TERMINAL_CACHE.clear()
def _explode(*_args, **_kwargs):
raise RuntimeError("network down")
@@ -837,14 +838,14 @@ def test_scan_terminal_endpoint_forwards_filters(monkeypatch):
def test_scan_terminal_cache_key_includes_filter_dimensions():
first = _scan_terminal_cache_key(
first = scan_terminal_cache_key(
{
"scan_mode": "tradable",
"time_range": "today",
"limit": 25,
}
)
second = _scan_terminal_cache_key(
second = scan_terminal_cache_key(
{
"scan_mode": "trend",
"time_range": "week",
+408
View File
@@ -0,0 +1,408 @@
from __future__ import annotations
import re
from datetime import datetime
from typing import Any, Dict, List, Optional
from web.scan_city_ai_helpers import _safe_float, _truncate_ai_text
def _compact_ai_candidate(row: Dict[str, Any]) -> Dict[str, Any]:
return {
"id": row.get("id"),
"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"),
"market_probability": row.get("market_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"),
"cluster_role": row.get("cluster_role"),
"model_cluster_sources": _compact_ai_model_sources(row),
"metar_context": row.get("metar_context") or {},
"window_phase": row.get("window_phase"),
"peak_window_label": row.get("peak_window_label"),
"minutes_until_peak_start": row.get("minutes_until_peak_start"),
"minutes_until_peak_end": row.get("minutes_until_peak_end"),
"trend_alignment": row.get("trend_alignment"),
"tradable": row.get("tradable"),
"accepting_orders": row.get("accepting_orders"),
}
def _normalize_ai_city_key(value: Any) -> str:
return str(value or "").strip().lower().replace(" ", "").replace("-", "").replace("_", "")
def _compact_ai_model_sources(row: Dict[str, Any]) -> List[Dict[str, Any]]:
raw_sources = row.get("model_cluster_sources")
if not isinstance(raw_sources, dict):
return []
sources: List[Dict[str, Any]] = []
for name, value in raw_sources.items():
if _safe_float(value) is None:
continue
sources.append({"model": str(name), "value": value})
return sources[:12]
def _observation_sort_key(point: Dict[str, Any]) -> tuple[int, str]:
raw_time = str(point.get("time") or "").strip()
try:
parsed = datetime.fromisoformat(raw_time.replace("Z", "+00:00"))
return parsed.hour * 60 + parsed.minute, raw_time
except Exception:
pass
match = re.search(r"(\d{1,2}):(\d{2})", raw_time)
if match:
hour = max(0, min(23, int(match.group(1))))
minute = max(0, min(59, int(match.group(2))))
return hour * 60 + minute, raw_time
return 9999, raw_time
def _compact_observation_points(raw_points: Any, limit: int = 24) -> List[Dict[str, Any]]:
if not isinstance(raw_points, list):
return []
points: List[Dict[str, Any]] = []
for item in raw_points:
if isinstance(item, dict):
temp = _safe_float(item.get("temp"))
time_value = str(item.get("time") or item.get("obs_time") or item.get("time_label") or "").strip()
elif isinstance(item, (list, tuple)) and len(item) >= 2:
time_value = str(item[0] or "").strip()
temp = _safe_float(item[1])
else:
continue
if temp is None or not time_value:
continue
points.append({"time": time_value, "temp": temp})
sorted_points = sorted(points, key=_observation_sort_key)
return sorted_points[-max(1, int(limit)) :]
def _compact_ai_text(value: Any, limit: int = 700) -> Optional[str]:
text = _truncate_ai_text(value, limit).strip()
return text or None
def _compact_hourly_context(raw_hourly: Any) -> Dict[str, Any]:
if not isinstance(raw_hourly, dict):
return {}
times = raw_hourly.get("times") or raw_hourly.get("time") or []
temps = raw_hourly.get("temps") or raw_hourly.get("temperature_2m") or []
radiation = raw_hourly.get("radiation") or raw_hourly.get("shortwave_radiation") or []
if not isinstance(times, list) or not isinstance(temps, list):
return {}
points: List[Dict[str, Any]] = []
for idx, raw_time in enumerate(times):
temp = _safe_float(temps[idx] if idx < len(temps) else None)
if temp is None:
continue
time_text = str(raw_time or "").strip()
if "T" in time_text:
time_text = time_text.split("T", 1)[1][:5]
elif len(time_text) > 5:
time_text = time_text[:5]
point: Dict[str, Any] = {"time": time_text, "temp": temp}
rad = _safe_float(radiation[idx] if isinstance(radiation, list) and idx < len(radiation) else None)
if rad is not None:
point["radiation"] = rad
points.append(point)
if not points:
return {}
max_point = max(points, key=lambda item: _safe_float(item.get("temp")) or -999.0)
sample_indexes = {
idx
for idx in range(len(points))
if idx % 2 == 0 or idx >= len(points) - 4 or points[idx] is max_point
}
samples = [points[idx] for idx in sorted(sample_indexes)][-14:]
return {
"sample_count": len(points),
"forecast_hourly_max": max_point,
"samples": samples,
}
def _compact_taf_context(raw_taf_data: Any) -> Dict[str, Any]:
if not isinstance(raw_taf_data, dict):
return {}
signal = raw_taf_data.get("signal") if isinstance(raw_taf_data.get("signal"), dict) else {}
source = signal or raw_taf_data
raw_taf = raw_taf_data.get("raw_taf") or source.get("raw_taf")
compact: Dict[str, Any] = {
"available": bool(source.get("available") or raw_taf),
"raw_taf": _compact_ai_text(raw_taf, 900),
"issue_time": raw_taf_data.get("issue_time") or source.get("issue_time"),
"valid_time_from": raw_taf_data.get("valid_time_from") or source.get("valid_time_from"),
"valid_time_to": raw_taf_data.get("valid_time_to") or source.get("valid_time_to"),
"peak_window": source.get("peak_window"),
"suppression_level": source.get("suppression_level"),
"disruption_level": source.get("disruption_level"),
"wind_shift": source.get("wind_shift"),
"wind_regimes": source.get("wind_regimes"),
"summary_zh": _compact_ai_text(source.get("summary_zh"), 260),
"summary_en": _compact_ai_text(source.get("summary_en"), 260),
}
segments = source.get("segments") if isinstance(source.get("segments"), list) else []
markers = source.get("markers") if isinstance(source.get("markers"), list) else []
if segments:
compact["segments"] = segments[:3]
if markers:
compact["markers"] = markers[:4]
return {key: value for key, value in compact.items() if value not in (None, "", [])}
def _compact_vertical_context(raw_vertical: Any) -> Dict[str, Any]:
if not isinstance(raw_vertical, dict):
return {}
keys = [
"source",
"window_start",
"window_end",
"suppression_risk",
"trigger_risk",
"mixing_strength",
"shear_risk",
"heating_setup",
"heating_score",
"summary_zh",
"summary_en",
]
compact: Dict[str, Any] = {}
for key in keys:
value = raw_vertical.get(key)
if isinstance(value, str):
value = _compact_ai_text(value, 280)
if value not in (None, "", []):
compact[key] = value
return compact
def _compact_intraday_context(raw_intraday: Any) -> Dict[str, Any]:
if not isinstance(raw_intraday, dict):
return {}
compact: Dict[str, Any] = {}
for key in [
"headline",
"headline_en",
"confidence",
"base_case_bucket",
"upside_bucket",
"downside_bucket",
"next_observation_time",
"peak_window",
]:
value = raw_intraday.get(key)
if isinstance(value, str):
value = _compact_ai_text(value, 220)
if value not in (None, "", []):
compact[key] = value
signals = raw_intraday.get("signal_contributions")
if isinstance(signals, list):
compact["signal_contributions"] = [
{
"label": item.get("label"),
"label_en": item.get("label_en"),
"direction": item.get("direction"),
"strength": item.get("strength"),
"summary": _compact_ai_text(item.get("summary"), 180),
"summary_en": _compact_ai_text(item.get("summary_en"), 180),
}
for item in signals[:4]
if isinstance(item, dict)
]
return compact
def _build_metar_decision_context(data: Dict[str, Any]) -> Dict[str, Any]:
today_obs = _compact_observation_points(data.get("metar_today_obs"), 36)
recent_obs = _compact_observation_points(data.get("metar_recent_obs"), 12)
settlement_obs = _compact_observation_points(data.get("settlement_today_obs"), 36)
airport_current = data.get("airport_current") if isinstance(data.get("airport_current"), dict) else {}
metar_status = data.get("metar_status") if isinstance(data.get("metar_status"), dict) else {}
source_obs = today_obs or recent_obs or settlement_obs
trend_source = recent_obs or source_obs[-4:]
last_point = source_obs[-1] if source_obs else {}
first_trend = trend_source[0] if trend_source else {}
last_trend = trend_source[-1] if trend_source else {}
max_point = None
for point in source_obs:
if max_point is None or float(point["temp"]) >= float(max_point["temp"]):
max_point = point
last_temp = _safe_float(last_point.get("temp"))
first_temp = _safe_float(first_trend.get("temp"))
trend_last_temp = _safe_float(last_trend.get("temp"))
trend_delta = (
trend_last_temp - first_temp
if trend_last_temp is not None and first_temp is not None and len(trend_source) >= 2
else None
)
station = data.get("risk") if isinstance(data.get("risk"), dict) else {}
current = data.get("current") if isinstance(data.get("current"), dict) else {}
settlement_station = data.get("settlement_station") if isinstance(data.get("settlement_station"), dict) else {}
settlement_source = str(
current.get("settlement_source")
or settlement_station.get("settlement_source")
or "metar"
).strip().lower()
is_hko = settlement_source == "hko"
source_label = "HKO" if is_hko else "METAR"
return {
"source": source_label,
"is_airport_metar": not is_hko,
"station": (
current.get("station_code")
or settlement_station.get("settlement_station_code")
or station.get("icao")
or airport_current.get("station_code")
),
"station_label": (
current.get("station_name")
or settlement_station.get("settlement_station_label")
or station.get("airport")
or airport_current.get("station_label")
),
"today_obs": today_obs[-12:],
"recent_obs": recent_obs[-8:],
"settlement_today_obs": settlement_obs[-12:],
"obs_count": len(source_obs),
"last_time": last_point.get("time"),
"last_temp": last_temp,
"max_temp": _safe_float((max_point or {}).get("temp")),
"max_time": (max_point or {}).get("time"),
"trend_delta": trend_delta,
"stale_for_today": bool(metar_status.get("stale_for_today")),
"available_for_today": bool(metar_status.get("available_for_today")),
"last_observation_time": metar_status.get("last_observation_time"),
"airport_current_temp": _safe_float(airport_current.get("temp")),
"airport_max_so_far": _safe_float(airport_current.get("max_so_far")),
"airport_obs_time": airport_current.get("obs_time"),
"airport_report_time": airport_current.get("report_time"),
"airport_raw_metar": airport_current.get("raw_metar"),
"airport_wx_desc": airport_current.get("wx_desc"),
"airport_cloud_desc": airport_current.get("cloud_desc"),
"airport_visibility_mi": _safe_float(airport_current.get("visibility_mi")),
"airport_wind_speed_kt": _safe_float(airport_current.get("wind_speed_kt")),
"airport_wind_dir": _safe_float(airport_current.get("wind_dir")),
"airport_humidity": _safe_float(airport_current.get("humidity")),
}
def _city_observation_anchor(data: Dict[str, Any]) -> Dict[str, Any]:
current = data.get("current") if isinstance(data.get("current"), dict) else {}
settlement_station = data.get("settlement_station") if isinstance(data.get("settlement_station"), dict) else {}
airport_current = data.get("airport_current") if isinstance(data.get("airport_current"), dict) else {}
risk = data.get("risk") if isinstance(data.get("risk"), dict) else {}
source = str(
current.get("settlement_source")
or settlement_station.get("settlement_source")
or "metar"
).strip().lower()
is_hko = source == "hko"
if is_hko:
station_code = (
current.get("station_code")
or settlement_station.get("settlement_station_code")
or "HKO"
)
station_label = (
current.get("station_name")
or settlement_station.get("settlement_station_label")
or "Hong Kong Observatory"
)
return {
"source": "hko",
"source_label": "Hong Kong Observatory",
"is_airport_metar": False,
"station_code": station_code,
"station_label": station_label,
"read_label_zh": "香港天文台观测解读",
"read_label_en": "Hong Kong Observatory observation read",
"instruction_zh": "该城市使用香港天文台/HKO 官方站点观测,不是机场 METAR;不得称为机场报文或 METAR。",
"instruction_en": "This city uses Hong Kong Observatory/HKO station observations, not an airport METAR; do not call it an airport bulletin or METAR.",
}
return {
"source": "metar",
"source_label": "METAR",
"is_airport_metar": True,
"station_code": risk.get("icao") or airport_current.get("station_code"),
"station_label": risk.get("airport") or airport_current.get("station_label"),
"read_label_zh": "机场报文解读",
"read_label_en": "airport-bulletin read",
"instruction_zh": "该城市使用机场 METAR/TAF 作为日内实况证据。",
"instruction_en": "This city uses airport METAR/TAF as intraday observation evidence.",
}
def _compact_ai_city_group(rows: List[Dict[str, Any]]) -> Dict[str, Any]:
first = rows[0]
return {
"city": first.get("city"),
"city_display_name": first.get("city_display_name") or first.get("display_name") or first.get("city"),
"selected_date": first.get("selected_date") or first.get("local_date"),
"local_time": first.get("local_time"),
"temp_symbol": first.get("temp_symbol") or first.get("target_unit"),
"current_temp": first.get("current_temp"),
"current_max_so_far": first.get("current_max_so_far"),
"deb_prediction": first.get("deb_prediction"),
"window_phase": first.get("window_phase"),
"remaining_window_minutes": first.get("remaining_window_minutes"),
"peak_window_label": first.get("peak_window_label"),
"minutes_until_peak_start": first.get("minutes_until_peak_start"),
"minutes_until_peak_end": first.get("minutes_until_peak_end"),
"metar_context": first.get("metar_context") or {},
"model_cluster": {
"core_low": first.get("cluster_core_low"),
"core_high": first.get("cluster_core_high"),
"median": first.get("cluster_median"),
"deb_reference": first.get("cluster_deb_reference"),
"model_count": first.get("cluster_model_count"),
"sources": _compact_ai_model_sources(first),
},
"contracts": [_compact_ai_candidate(row) for row in rows],
}
def build_scan_ai_prompt(payload: Dict[str, Any], *, max_rows: int) -> Dict[str, Any]:
raw_rows = [
row
for row in (payload.get("rows") or [])[:max_rows]
if isinstance(row, dict) and row.get("id")
]
grouped: Dict[str, List[Dict[str, Any]]] = {}
for row in raw_rows:
key = "|".join(
[
_normalize_ai_city_key(row.get("city") or row.get("city_display_name")),
str(row.get("selected_date") or row.get("local_date") or ""),
]
)
grouped.setdefault(key, []).append(row)
cities = [_compact_ai_city_group(rows) for rows in grouped.values() if rows]
sent_contracts = sum(len(city.get("contracts") or []) for city in cities)
return {
"schema_version": "city_forecast_v1",
"snapshot_id": payload.get("snapshot_id"),
"generated_at": payload.get("generated_at"),
"summary": payload.get("summary") or {},
"filters": payload.get("filters") or {},
"city_count": len(cities),
"candidate_row_count": len(raw_rows),
"cities": cities,
"_polyweather_input_meta": {
"sent_cities": len(cities),
"sent_contracts": sent_contracts,
},
}
+255
View File
@@ -0,0 +1,255 @@
from __future__ import annotations
from datetime import datetime
from typing import Any, Dict, List, Optional
from web.scan_city_ai_helpers import _safe_float
from web.scan_terminal_ai_compact import _normalize_ai_city_key
from web.scan_terminal_filters import safe_int as _safe_int
from web.scan_terminal_metar_gate import _apply_metar_gate_to_row
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 _normalize_ai_city_theses(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 not isinstance(item, dict):
continue
city = str(item.get("city") or item.get("city_name") or "").strip()
if not city:
continue
out.append({**item, "city": city})
return out
def _normalize_ai_city_forecasts(ai_raw: Dict[str, Any]) -> List[Dict[str, Any]]:
raw_items = (
ai_raw.get("city_forecasts")
or ai_raw.get("city_predictions")
or ai_raw.get("city_max_forecasts")
or ai_raw.get("city_theses")
)
if not isinstance(raw_items, list):
return []
out: List[Dict[str, Any]] = []
for item in raw_items:
if not isinstance(item, dict):
continue
city = str(item.get("city") or item.get("city_name") or "").strip()
if not city:
continue
predicted = (
item.get("predicted_max")
if item.get("predicted_max") is not None
else item.get("max_temp")
if item.get("max_temp") is not None
else item.get("prediction")
)
out.append(
{
**item,
"city": city,
"predicted_max": predicted,
"range_low": item.get("range_low") if item.get("range_low") is not None else item.get("low"),
"range_high": item.get("range_high") if item.get("range_high") is not None else item.get("high"),
"reasoning_zh": item.get("reasoning_zh") or item.get("thesis_zh") or item.get("summary_zh"),
"reasoning_en": item.get("reasoning_en") or item.get("thesis_en") or item.get("summary_en"),
}
)
return out
def merge_scan_ai_result(
payload: Dict[str, Any],
ai_raw: Dict[str, Any],
*,
model: str,
max_rows: int,
timeout_sec: int,
cache_ttl_sec: int,
base_url: str,
cached: bool = False,
duration_ms: Optional[int] = None,
input_rows: Optional[int] = None,
) -> 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"))
watchlist = _normalize_ai_items(ai_raw.get("watchlist"))
city_theses = _normalize_ai_city_theses(ai_raw.get("city_theses"))
city_forecasts = _normalize_ai_city_forecasts(ai_raw)
contract_notes = _normalize_ai_items(ai_raw.get("contract_notes"))
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()
watchlist_ids = {str(item.get("row_id")) for item in watchlist}
thesis_by_city: Dict[str, Dict[str, Any]] = {}
for item in city_theses:
key = _normalize_ai_city_key(item.get("city"))
if key:
thesis_by_city[key] = item
forecast_by_city: Dict[str, Dict[str, Any]] = {}
for item in city_forecasts:
key = _normalize_ai_city_key(item.get("city"))
if key:
forecast_by_city[key] = item
for row in rows:
city_key = _normalize_ai_city_key(row.get("city"))
display_key = _normalize_ai_city_key(row.get("city_display_name"))
thesis = thesis_by_city.get(city_key) or thesis_by_city.get(display_key)
forecast = forecast_by_city.get(city_key) or forecast_by_city.get(display_key)
if thesis:
row["ai_city_thesis_zh"] = thesis.get("thesis_zh") or thesis.get("summary_zh")
row["ai_city_thesis_en"] = thesis.get("thesis_en") or thesis.get("summary_en")
row["ai_city_confidence"] = thesis.get("confidence")
row["ai_city_model_cluster_note"] = thesis.get("model_cluster_note")
if forecast:
row["ai_predicted_max"] = _safe_float(forecast.get("predicted_max"))
row["ai_predicted_low"] = _safe_float(forecast.get("range_low"))
row["ai_predicted_high"] = _safe_float(forecast.get("range_high"))
row["ai_forecast_unit"] = forecast.get("unit") or row.get("temp_symbol")
row["ai_forecast_confidence"] = forecast.get("confidence")
row["ai_peak_window_zh"] = forecast.get("peak_window_zh")
row["ai_peak_window_en"] = forecast.get("peak_window_en")
row["ai_airport_metar_read_zh"] = forecast.get("metar_read_zh")
row["ai_airport_metar_read_en"] = forecast.get("metar_read_en")
row["ai_forecast_reason_zh"] = forecast.get("reasoning_zh")
row["ai_forecast_reason_en"] = forecast.get("reasoning_en")
row["ai_city_model_cluster_note"] = forecast.get("model_cluster_note") or row.get("ai_city_model_cluster_note")
row["ai_city_thesis_zh"] = row.get("ai_city_thesis_zh") or forecast.get("reasoning_zh")
row["ai_city_thesis_en"] = row.get("ai_city_thesis_en") or forecast.get("reasoning_en")
for item in contract_notes:
row = by_id.get(str(item.get("row_id")))
if not row:
continue
row["ai_forecast_match"] = item.get("forecast_match") or item.get("match")
row["ai_forecast_match_reason_zh"] = item.get("reason_zh") or item.get("reason")
row["ai_forecast_match_reason_en"] = item.get("reason_en")
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 item in watchlist:
row = by_id.get(str(item.get("row_id")))
if not row:
continue
row["ai_watchlist_reason_zh"] = item.get("reason_zh") or item.get("reason")
row["ai_watchlist_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"
if row_id in watchlist_ids and row.get("ai_decision") == "neutral":
row["ai_decision"] = "watchlist"
_apply_metar_gate_to_row(row)
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,
)
input_meta = ai_raw.get("_polyweather_input_meta")
sent_cities = input_meta.get("sent_cities") if isinstance(input_meta, dict) else None
sent_contracts = input_meta.get("sent_contracts") if isinstance(input_meta, dict) else None
ai_scan = {
"status": "ready",
"stage": "completed",
"model": model,
"cached": cached,
"generated_at": datetime.utcnow().isoformat() + "Z",
"snapshot_id": payload.get("snapshot_id"),
"input_rows": input_rows if input_rows is not None else len(payload.get("rows") or []),
"sent_rows": sent_contracts if sent_contracts is not None else min(len(payload.get("rows") or []), max_rows),
"sent_cities": sent_cities,
"sent_contracts": sent_contracts,
"duration_ms": duration_ms,
"timeout_sec": timeout_sec,
"cache_ttl_sec": cache_ttl_sec,
"provider": "deepseek",
"base_url": base_url,
"summary_zh": ai_raw.get("summary_zh"),
"summary_en": ai_raw.get("summary_en"),
"city_forecasts": city_forecasts,
"contract_notes": contract_notes,
"city_theses": city_theses,
"watchlist": watchlist,
"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"),
"watchlist_count": sum(1 for row in rows if row.get("ai_decision") == "watchlist"),
}
meta = ai_raw.get("_polyweather_meta")
if isinstance(meta, dict):
ai_scan["usage"] = meta.get("usage")
ai_scan["finish_reason"] = meta.get("finish_reason")
merged = {
**payload,
"rows": rows,
"top_signal": top_signal,
"ai_scan": ai_scan,
}
return merged
+163
View File
@@ -0,0 +1,163 @@
from __future__ import annotations
import hashlib
import json
import threading
import time
from datetime import datetime
from typing import Any, Dict, Optional
_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]] = {}
def scan_terminal_cache_key(filters: Dict[str, Any]) -> str:
return json.dumps(filters, ensure_ascii=True, sort_keys=True)
def get_cached_scan_terminal_payload(
filters: Dict[str, Any],
*,
ttl_sec: int,
) -> 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(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 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 {}
now = time.time()
with _SCAN_TERMINAL_CACHE_LOCK:
_SCAN_TERMINAL_CACHE[cache_key] = {
"t": now,
"payload": dict(payload),
"success_t": now,
"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 mark_scan_terminal_refreshing(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)
return True
def clear_scan_terminal_refreshing(filters: Dict[str, Any]) -> None:
cache_key = scan_terminal_cache_key(filters)
with _SCAN_TERMINAL_CACHE_LOCK:
_SCAN_TERMINAL_REFRESHING.discard(cache_key)
def scan_ai_cache_key(
snapshot_id: str,
filters: Dict[str, Any],
*,
max_rows: int,
model: str,
) -> str:
raw = json.dumps(
{
"schema_version": "city_forecast_v1",
"snapshot_id": snapshot_id,
"filters": filters,
"model": model,
"max_rows": 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],
*,
max_rows: int,
model: str,
ttl_sec: int,
) -> Optional[Dict[str, Any]]:
cache_key = scan_ai_cache_key(
snapshot_id,
filters,
max_rows=max_rows,
model=model,
)
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(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],
*,
max_rows: int,
model: str,
) -> None:
cache_key = scan_ai_cache_key(
snapshot_id,
filters,
max_rows=max_rows,
model=model,
)
with _SCAN_TERMINAL_AI_CACHE_LOCK:
_SCAN_TERMINAL_AI_CACHE[cache_key] = {
"cached_at": time.time(),
"result": result,
}
+182
View File
@@ -0,0 +1,182 @@
from __future__ import annotations
from datetime import datetime, timedelta
from typing import Any, Dict, List
from web.core import CITIES
from web.analysis_service import _analyze, _build_city_market_scan_payload
from web.scan_city_ai_helpers import _safe_float
from web.scan_terminal_ai_compact import _build_metar_decision_context
from web.scan_terminal_filters import (
market_region_from_tz_offset as _market_region_from_tz_offset,
safe_int as _safe_int,
)
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 {}
city_meta = CITIES.get(city) or {}
tz_offset = _safe_int(city_meta.get("tz"), 0)
market_region = _market_region_from_tz_offset(tz_offset)
metar_context = _build_metar_decision_context(data)
return {
**row,
"id": str(row.get("id") or f"{city}|{selected_date}|{market_slug}|{side}"),
"city": city,
"city_display_name": display_name,
"trading_region": market_region["key"],
"trading_region_label": market_region["label_en"],
"trading_region_label_zh": market_region["label_zh"],
"tz_offset_seconds": tz_offset,
"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"),
"metar_context": metar_context,
"metar_today_obs": metar_context.get("today_obs") or [],
"metar_recent_obs": metar_context.get("recent_obs") or [],
"settlement_today_obs": metar_context.get("settlement_today_obs") or [],
"metar_status": {
"available_for_today": metar_context.get("available_for_today"),
"stale_for_today": metar_context.get("stale_for_today"),
"last_observation_time": metar_context.get("last_observation_time"),
"last_temp": metar_context.get("last_temp"),
},
"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 [],
"distribution_full": scan.get("distribution_full") or scan.get("distribution_preview") or row.get("distribution_preview") or [],
"model_cluster_sources": daily_entry.get("models") if isinstance(daily_entry.get("models"), dict) else data.get("multi_model"),
"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,
}
+72
View File
@@ -0,0 +1,72 @@
from __future__ import annotations
from typing import Any, Dict, Optional
from web.scan_city_ai_helpers import _safe_float
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 market_region_from_tz_offset(tz_offset_seconds: Any) -> Dict[str, str]:
tz_offset = safe_int(tz_offset_seconds, 0)
if tz_offset <= -7200:
return {
"key": "americas",
"label_en": "Americas",
"label_zh": "美洲",
}
if tz_offset >= 14400:
return {
"key": "asia_pacific",
"label_en": "Asia-Pacific",
"label_zh": "亚太",
}
return {
"key": "europe_africa",
"label_en": "Europe / Africa",
"label_zh": "欧洲 / 非洲",
}
+154
View File
@@ -0,0 +1,154 @@
from __future__ import annotations
import re
from typing import Any, Dict, Optional
from web.scan_city_ai_helpers import _safe_float
from web.scan_terminal_filters import safe_int as _safe_int
def _target_range_from_row(row: Dict[str, Any]) -> tuple[Optional[float], Optional[float]]:
lower = _safe_float(row.get("target_lower"))
upper = _safe_float(row.get("target_upper"))
if lower is not None or upper is not None:
return lower, upper
threshold = _safe_float(row.get("target_threshold"))
target_value = _safe_float(row.get("target_value"))
raw_label = str(row.get("target_label") or row.get("action") or "")
numbers = [float(match.group(0)) for match in re.finditer(r"-?\d+(?:\.\d+)?", raw_label)]
if len(numbers) >= 2:
return min(numbers[0], numbers[1]), max(numbers[0], numbers[1])
value = threshold if threshold is not None else target_value if target_value is not None else (numbers[0] if numbers else None)
if value is None:
return None, None
if re.search(r"(\+|above|higher|or\s+higher|>=|≥|以上)", raw_label, re.I):
return value, None
if re.search(r"(below|or\s+below|<=|≤|以下)", raw_label, re.I):
return None, value
return value, value
def _metar_gate_for_row(row: Dict[str, Any]) -> Optional[Dict[str, Any]]:
context = row.get("metar_context") if isinstance(row.get("metar_context"), dict) else {}
side = str(row.get("side") or "").strip().lower()
if side not in {"yes", "no"}:
return None
obs_count = _safe_int(context.get("obs_count"), 0)
if obs_count <= 0 or context.get("stale_for_today"):
return {
"decision": "downgrade",
"reason_zh": "V4 未拿到同日 METAR 实测,不能只凭 edge/Kelly 给出交易。",
"reason_en": "V4 has no same-day METAR observations, so edge/Kelly alone cannot drive a trade.",
}
lower, upper = _target_range_from_row(row)
max_temp = _safe_float(context.get("max_temp"))
last_temp = _safe_float(context.get("last_temp"))
trend_delta = _safe_float(context.get("trend_delta"))
if max_temp is None or (lower is None and upper is None):
return None
unit = str(row.get("target_unit") or row.get("temp_symbol") or "")
epsilon = 0.7 if "F" in unit.upper() else 0.4
phase = str(row.get("window_phase") or "").lower()
remaining = _safe_float(row.get("remaining_window_minutes"))
minutes_until_peak_start = _safe_float(row.get("minutes_until_peak_start"))
is_late = phase in {"active_peak", "post_peak"} or (remaining is not None and remaining <= 180)
is_before_peak = phase in {"early_today", "setup_today", "tomorrow", "week_ahead"} or (
minutes_until_peak_start is not None and minutes_until_peak_start > 0
)
is_falling = trend_delta is not None and trend_delta <= -epsilon
is_not_rising = trend_delta is not None and trend_delta <= epsilon
above_upper = upper is not None and max_temp > upper + epsilon
below_lower = lower is not None and max_temp < lower - epsilon
inside_bucket = (
(lower is None or max_temp >= lower - epsilon)
and (upper is None or max_temp <= upper + epsilon)
)
if side == "no":
if above_upper:
return {
"decision": "approve",
"reason_zh": "METAR 实测最高已越过目标桶上沿,V4 确认 BUY NO 有实测支撑。",
"reason_en": "METAR max has already moved above the bucket, so V4 confirms BUY NO has observation support.",
}
if below_lower and (is_late or is_falling or is_not_rising):
if is_before_peak and not is_late:
return {
"decision": "watchlist",
"reason_zh": "峰值窗口尚未到来,METAR 暂未触达不能直接确认 BUY NO,V4 先列观察。",
"reason_en": "The peak window has not arrived, so a still-low METAR path cannot confirm BUY NO yet; V4 keeps it on watch.",
}
return {
"decision": "approve",
"reason_zh": "METAR 最高仍低于目标桶且近期走势不强,V4 确认 BUY NO 优先。",
"reason_en": "METAR max remains below the bucket and recent observations are not strengthening, so V4 favors BUY NO.",
}
if inside_bucket and is_late and is_not_rising:
return {
"decision": "downgrade",
"reason_zh": "METAR 最高仍贴近目标桶,V4 不允许只因 edge 高就直接交易 NO。",
"reason_en": "METAR max is still close to the target bucket, so V4 will not trade NO on edge alone.",
}
else:
if above_upper:
return {
"decision": "veto",
"reason_zh": "METAR 实测最高已越过目标桶上沿,V4 排除该 BUY YES。",
"reason_en": "METAR max has already exceeded the bucket, so V4 vetoes this BUY YES.",
}
if below_lower and (is_late or is_falling or is_not_rising):
if is_before_peak and not is_late:
return {
"decision": "watchlist",
"reason_zh": "峰值窗口尚未到来,METAR 未触达目标桶只能说明仍需等待峰值验证,V4 暂列观察。",
"reason_en": "The peak window has not arrived, so METAR not reaching the bucket only means the setup still needs peak-window confirmation; V4 keeps it on watch.",
}
return {
"decision": "downgrade",
"reason_zh": "METAR 最高仍未触达目标桶且走势不强,V4 将 BUY YES 降级观察。",
"reason_en": "METAR max has not reached the bucket and recent observations are weak, so V4 downgrades BUY YES.",
}
if inside_bucket:
return {
"decision": "approve",
"reason_zh": "METAR 实测最高已落入目标桶,V4 认为 BUY YES 有实测依据,但仍需防止继续升穿上沿。",
"reason_en": "METAR max is inside the target bucket, so V4 sees observation support for BUY YES while monitoring an overshoot.",
}
if last_temp is not None and trend_delta is not None:
direction = "走弱" if trend_delta < -epsilon else "走强" if trend_delta > epsilon else "横盘"
return {
"decision": "watchlist",
"reason_zh": f"METAR 最新 {last_temp:.1f},近期{direction}V4 暂不把该合约升级为最终交易。",
"reason_en": f"Latest METAR is {last_temp:.1f} with a recent {'downtrend' if trend_delta < -epsilon else 'uptrend' if trend_delta > epsilon else 'flat trend'}, so V4 keeps this as watchlist.",
}
return None
def _apply_metar_gate_to_row(row: Dict[str, Any]) -> None:
gate = _metar_gate_for_row(row)
if not gate:
return
decision = str(gate.get("decision") or "").lower()
row["v4_metar_decision"] = decision
row["v4_metar_reason_zh"] = gate.get("reason_zh")
row["v4_metar_reason_en"] = gate.get("reason_en")
current_decision = str(row.get("ai_decision") or "").lower()
hard_decisions = {"veto", "downgrade"}
if decision == "veto":
row["ai_decision"] = "veto"
row.pop("ai_rank", None)
elif decision == "downgrade" and current_decision != "veto":
row["ai_decision"] = "downgrade"
row.pop("ai_rank", None)
elif decision == "approve" and current_decision not in hard_decisions:
row["ai_decision"] = "approve"
elif decision == "watchlist" and current_decision not in {"approve", "veto", "downgrade"}:
row["ai_decision"] = "watchlist"
if decision in {"approve", "veto", "downgrade"}:
row["ai_reason_zh"] = gate.get("reason_zh") or row.get("ai_reason_zh")
row["ai_reason_en"] = gate.get("reason_en") or row.get("ai_reason_en")
+86
View File
@@ -0,0 +1,86 @@
from __future__ import annotations
import hashlib
import json
from datetime import datetime
from typing import Any, Dict, List, Optional
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 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": [],
}
+84
View File
@@ -0,0 +1,84 @@
from __future__ import annotations
from typing import Any, Dict, List, Optional
from web.scan_city_ai_helpers import _safe_float
def build_ranked_scan_terminal_result(
*,
city_results: List[Dict[str, Any]],
filters: Dict[str, Any],
total_city_count: int,
failed_city_count: int,
) -> Dict[str, Any]:
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]] = [
{
**row,
"rank": index,
}
for index, row in enumerate(primary_rows[: filters["limit"]], start=1)
]
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: Optional[float] = 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: Optional[float] = 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": total_city_count,
"scanned_city_count": len(city_results),
"failed_city_count": failed_city_count,
}
return {
"primary_rows": primary_rows,
"ranked_rows": ranked_rows,
"summary": summary,
"top_signal": top_signal,
}
File diff suppressed because it is too large Load Diff