diff --git a/frontend/components/dashboard/scan-terminal/MarketOverviewBanner.module.css b/frontend/components/dashboard/scan-terminal/MarketOverviewBanner.module.css
new file mode 100644
index 00000000..3de641b0
--- /dev/null
+++ b/frontend/components/dashboard/scan-terminal/MarketOverviewBanner.module.css
@@ -0,0 +1,192 @@
+.root {
+ display: flex;
+ flex-direction: column;
+ border: 1px solid var(--color-border-default, rgba(159, 178, 199, 0.16));
+ border-radius: var(--radius-md, 10px);
+ background: var(--color-bg-card, rgba(17, 26, 46, 0.88));
+ backdrop-filter: var(--glass-blur-1, blur(10px));
+ margin: 0 0 var(--space-3, 12px) 0;
+ transition: background 0.2s;
+}
+
+.root.loading {
+ opacity: 0.7;
+ flex-direction: row;
+ align-items: center;
+ gap: 8px;
+ padding: 10px 14px;
+ font-size: 13px;
+ color: var(--color-text-secondary, #9FB2C7);
+}
+
+.icon {
+ color: var(--color-accent-primary, #4DA3FF);
+ flex-shrink: 0;
+}
+
+.header {
+ display: flex;
+ align-items: center;
+ gap: 8px;
+ padding: 10px 14px;
+ cursor: pointer;
+ border: none;
+ background: transparent;
+ color: var(--color-text-primary, #E6EDF3);
+ font-size: 13px;
+ font-family: inherit;
+ text-align: left;
+ width: 100%;
+}
+
+.header:hover {
+ background: rgba(77, 163, 255, 0.04);
+}
+
+.badge {
+ display: inline-flex;
+ align-items: center;
+ gap: 5px;
+ color: var(--color-accent-primary, #4DA3FF);
+ font-weight: 700;
+ font-size: 12px;
+ letter-spacing: 0.03em;
+ flex-shrink: 0;
+}
+
+.preview {
+ flex: 1;
+ min-width: 0;
+ overflow: hidden;
+ text-overflow: ellipsis;
+ white-space: nowrap;
+ color: var(--color-text-secondary, #9FB2C7);
+ font-size: 12px;
+}
+
+.toggle {
+ flex-shrink: 0;
+ color: var(--color-text-muted, #7D8FA3);
+}
+
+.body {
+ padding: 0 14px 14px;
+ display: flex;
+ flex-direction: column;
+ gap: 10px;
+}
+
+.summary {
+ margin: 0;
+ font-size: 13px;
+ line-height: 1.65;
+ color: var(--color-text-primary, #E6EDF3);
+}
+
+.highlights {
+ list-style: none;
+ margin: 0;
+ padding: 0;
+ display: flex;
+ flex-wrap: wrap;
+ gap: 8px;
+}
+
+.highlightItem {
+ display: inline-flex;
+ align-items: baseline;
+ gap: 5px;
+ padding: 5px 10px;
+ border-radius: var(--radius-sm, 6px);
+ background: rgba(77, 163, 255, 0.08);
+ border: 1px solid rgba(77, 163, 255, 0.14);
+ font-size: 12px;
+ line-height: 1.5;
+ color: var(--color-text-secondary, #9FB2C7);
+ max-width: 100%;
+}
+
+.highlightItem strong {
+ color: var(--color-accent-primary, #4DA3FF);
+ font-weight: 700;
+ flex-shrink: 0;
+}
+
+.time {
+ font-size: 11px;
+ color: var(--color-text-muted, #7D8FA3);
+}
+
+/* Mobile < 768px */
+@media (max-width: 768px) {
+ .header {
+ padding: 8px 10px;
+ font-size: 12px;
+ }
+
+ .body {
+ padding: 0 10px 10px;
+ }
+
+ .summary {
+ font-size: 12px;
+ }
+
+ .highlights {
+ flex-direction: column;
+ gap: 6px;
+ }
+
+ .highlightItem {
+ font-size: 11px;
+ padding: 4px 8px;
+ }
+
+ .badge {
+ font-size: 11px;
+ }
+}
+
+/* Mobile < 640px */
+@media (max-width: 640px) {
+ .header {
+ padding: 6px 8px;
+ }
+
+ .preview {
+ display: none;
+ }
+
+ .body {
+ padding: 0 8px 8px;
+ }
+
+ .summary {
+ font-size: 11px;
+ line-height: 1.55;
+ }
+}
+
+/* Light theme */
+:global(.scan-terminal.light) .root {
+ background: rgba(255, 255, 255, 0.78);
+ border-color: rgba(15, 23, 42, 0.1);
+}
+
+:global(.scan-terminal.light) .summary {
+ color: #0F172A;
+}
+
+:global(.scan-terminal.light) .preview,
+:global(.scan-terminal.light) .highlightItem {
+ color: #475569;
+}
+
+:global(.scan-terminal.light) .highlightItem {
+ background: rgba(77, 163, 255, 0.06);
+ border-color: rgba(77, 163, 255, 0.18);
+}
+
+:global(.scan-terminal.light) .header:hover {
+ background: rgba(77, 163, 255, 0.05);
+}
diff --git a/frontend/components/dashboard/scan-terminal/MarketOverviewBanner.tsx b/frontend/components/dashboard/scan-terminal/MarketOverviewBanner.tsx
new file mode 100644
index 00000000..527c56cd
--- /dev/null
+++ b/frontend/components/dashboard/scan-terminal/MarketOverviewBanner.tsx
@@ -0,0 +1,132 @@
+"use client";
+
+import clsx from "clsx";
+import { ChevronDown, ChevronUp, Sparkles } from "lucide-react";
+import { useCallback, useEffect, useRef, useState } from "react";
+import styles from "./MarketOverviewBanner.module.css";
+import type { ScanOpportunityRow } from "@/lib/dashboard-types";
+
+interface OverviewPayload {
+ overview_zh: string;
+ overview_en: string;
+ highlights: Array<{ city: string; note_zh: string; note_en: string }>;
+ generated_at: string | null;
+}
+
+export function MarketOverviewBanner({
+ isEn,
+ isPro,
+ rows,
+}: {
+ isEn: boolean;
+ isPro: boolean;
+ rows: ScanOpportunityRow[];
+}) {
+ const [collapsed, setCollapsed] = useState(true);
+ const [data, setData] = useState
(null);
+ const [loading, setLoading] = useState(false);
+ const [error, setError] = useState(false);
+ const fetchedRef = useRef(false);
+
+ const fetchOverview = useCallback(async () => {
+ if (!isPro || !rows.length || fetchedRef.current) return;
+ fetchedRef.current = true;
+ setLoading(true);
+ setError(false);
+ try {
+ const resp = await fetch("/api/scan/terminal/overview", {
+ method: "POST",
+ headers: { "Content-Type": "application/json" },
+ body: JSON.stringify({
+ rows: rows.slice(0, 40).map((row) => ({
+ city: row.city ?? row.display_name ?? "",
+ display_name: row.display_name ?? row.city ?? "",
+ local_date: row.local_date ?? "",
+ deb_prediction: row.deb_prediction ?? null,
+ current_temp: row.current_temp ?? null,
+ current_max_so_far: row.current_max_so_far ?? null,
+ risk_level: row.risk_level ?? "",
+ temp_symbol: row.temp_symbol ?? "°C",
+ })),
+ }),
+ });
+ if (resp.ok) {
+ const json = await resp.json();
+ setData(json);
+ } else {
+ setError(true);
+ }
+ } catch {
+ setError(true);
+ } finally {
+ setLoading(false);
+ }
+ }, [isPro, rows]);
+
+ useEffect(() => {
+ if (isPro && rows.length > 0 && !fetchedRef.current) {
+ fetchOverview();
+ }
+ }, [isPro, rows.length, fetchOverview]);
+
+ if (!isPro || rows.length === 0) return null;
+ if (loading && !data) {
+ return (
+
+
+ {isEn ? "AI is generating market overview…" : "AI 正在生成市场概览…"}
+
+ );
+ }
+ if (error && !data) return null;
+
+ const overviewText = data ? (isEn ? data.overview_en : data.overview_zh) : "";
+ const highlights = data?.highlights ?? [];
+
+ return (
+
+
+
+ {!collapsed && (
+
+
{overviewText}
+ {highlights.length > 0 && (
+
+ {highlights.map((h) => (
+ -
+ {h.city}
+ {isEn ? h.note_en : h.note_zh}
+
+ ))}
+
+ )}
+ {data?.generated_at && (
+
+ )}
+
+ )}
+
+ );
+}
diff --git a/tests/test_web_observability.py b/tests/test_web_observability.py
index e5a18c83..426550cc 100644
--- a/tests/test_web_observability.py
+++ b/tests/test_web_observability.py
@@ -327,34 +327,30 @@ def test_city_ai_fallback_treats_stale_metar_as_background_not_anchor():
def test_city_ai_cache_key_changes_when_observation_fingerprint_changes():
+ """METAR 原文不变 → 缓存 key 不变(命中);METAR 原文变了 → key 变化(miss)。"""
base_input = {
- "schema_version": "single_city_forecast_v2",
"city": "Manila",
"local_date": "2026-04-28",
- "deb": {"prediction": 34.0},
"observation_anchor": {
"source": "METAR",
"is_airport_metar": True,
"station_code": "RPLL",
},
"airport_current": {
- "station_code": "RPLL",
- "temp": 34.0,
- "report_time": "03:00Z",
- "receipt_time": "2026-04-28T03:02:00Z",
+ "obs_time": "03:00Z",
"raw_metar": "RPLL 280300Z 34004KT CAVOK 34/24 Q1009",
},
"metar_context": {
"stale_for_today": False,
"last_observation_time": "03:00Z",
},
- "metar_today_obs": [{"time": "03:00Z", "temp": 34.0}],
}
changed_input = {
**base_input,
"airport_current": {
**base_input["airport_current"],
- "receipt_time": "2026-04-28T03:30:00Z",
+ "raw_metar": "RPLL 280330Z 36006KT 9999 FEW020 33/25 Q1010",
+ "obs_time": "03:30Z",
},
}
@@ -385,9 +381,10 @@ def test_city_ai_stream_request_only_asks_provider_for_observation_read():
user_payload = request_payload["messages"][1]["content"]
assert request_payload["stream"] is True
- assert request_payload["max_tokens"] <= 900
+ assert request_payload["max_tokens"] <= 1200
assert request_payload["max_tokens"] < scan_terminal_service.SCAN_AI_MAX_TOKENS
- assert '"required_json_keys": ["metar_read_zh", "metar_read_en", "predicted_max", "range_low", "range_high", "unit", "confidence", "final_judgment_zh", "final_judgment_en", "reasoning_zh", "reasoning_en"]' in user_payload
+ assert "taf_read_zh" in user_payload
+ assert "probability_read_zh" in user_payload
assert "predicted_max" in user_payload
assert "final_judgment" in user_payload
diff --git a/web/routers/scan.py b/web/routers/scan.py
index 968d98b6..de1b47be 100644
--- a/web/routers/scan.py
+++ b/web/routers/scan.py
@@ -6,6 +6,7 @@ from web.services.scan_api import (
get_scan_city_ai_forecast_payload,
get_scan_city_ai_stream_response,
get_scan_terminal_ai_payload,
+ get_scan_terminal_overview_payload,
get_scan_terminal_payload,
)
@@ -54,3 +55,8 @@ async def scan_terminal_ai_city(request: Request):
@router.post("/api/scan/terminal/ai-city/stream")
async def scan_terminal_ai_city_stream(request: Request):
return await get_scan_city_ai_stream_response(request)
+
+
+@router.post("/api/scan/terminal/overview")
+async def scan_terminal_overview(request: Request):
+ return await get_scan_terminal_overview_payload(request)
diff --git a/web/scan_city_ai_fallback.py b/web/scan_city_ai_fallback.py
index d37f9259..93abab7c 100644
--- a/web/scan_city_ai_fallback.py
+++ b/web/scan_city_ai_fallback.py
@@ -333,6 +333,34 @@ def _build_city_ai_fallback(
"range_low": range_low,
"range_high": range_high,
}
+ taf_data = ai_input.get("taf") if isinstance(ai_input.get("taf"), dict) else {}
+ if partial_ai.get("taf_read_zh") or partial_ai.get("taf_read_en"):
+ taf_zh = str(partial_ai.get("taf_read_zh") or partial_ai.get("taf_read_en") or "").strip()
+ taf_en = str(partial_ai.get("taf_read_en") or partial_ai.get("taf_read_zh") or "").strip()
+ elif taf_data.get("raw_taf"):
+ taf_zh = f"TAF 可用,需人工判读:{str(taf_data.get('raw_taf', ''))[:120]}"
+ taf_en = f"TAF available, manual read: {str(taf_data.get('raw_taf', ''))[:120]}"
+ else:
+ taf_zh = "无可用 TAF"
+ taf_en = "No TAF available"
+
+ prob_data = ai_input.get("probability") if isinstance(ai_input.get("probability"), dict) else {}
+ if partial_ai.get("probability_read_zh") or partial_ai.get("probability_read_en"):
+ prob_zh = str(partial_ai.get("probability_read_zh") or partial_ai.get("probability_read_en") or "").strip()
+ prob_en = str(partial_ai.get("probability_read_en") or partial_ai.get("probability_read_zh") or "").strip()
+ elif prob_data.get("top_buckets"):
+ top = prob_data["top_buckets"][0] if isinstance(prob_data["top_buckets"], list) else {}
+ if isinstance(top, dict) and top.get("label"):
+ skew_text = f",分布{'偏右' if prob_data.get('skew') == 'right' else '偏左' if prob_data.get('skew') == 'left' else '对称'}" if prob_data.get("skew") else ""
+ prob_zh = f"最高概率桶 {top.get('label', '')}({top.get('prob', '?')}%){skew_text}"
+ prob_en = f"Peak bucket {top.get('label', '')} ({top.get('prob', '?')}%){skew_text.replace('偏右', ', right-skewed').replace('偏左', ', left-skewed').replace('对称', ', symmetric')}"
+ else:
+ prob_zh = "概率分布数据可用但格式异常"
+ prob_en = "Probability data available but malformed"
+ else:
+ prob_zh = "概率分布暂未生成"
+ prob_en = "Probability distribution not yet generated"
+
return {
"predicted_max": partial_ai.get("predicted_max", predicted),
"range_low": partial_ai.get("range_low", range_low),
@@ -343,6 +371,10 @@ def _build_city_ai_fallback(
"final_judgment_en": final_en,
"metar_read_zh": metar_zh,
"metar_read_en": metar_en,
+ "taf_read_zh": taf_zh,
+ "taf_read_en": taf_en,
+ "probability_read_zh": prob_zh,
+ "probability_read_en": prob_en,
"reasoning_zh": reasoning_zh,
"reasoning_en": reasoning_en,
"risks_zh": risks_zh,
diff --git a/web/scan_city_ai_helpers.py b/web/scan_city_ai_helpers.py
index 953f3554..bfc44c97 100644
--- a/web/scan_city_ai_helpers.py
+++ b/web/scan_city_ai_helpers.py
@@ -7,6 +7,10 @@ from typing import Any, Dict, List, Optional
CITY_AI_REQUIRED_FIELDS = [
"metar_read_zh",
"metar_read_en",
+ "taf_read_zh",
+ "taf_read_en",
+ "probability_read_zh",
+ "probability_read_en",
"final_judgment_zh",
"final_judgment_en",
"predicted_max",
@@ -25,6 +29,10 @@ CITY_AI_REQUIRED_FIELDS = [
CITY_AI_STREAM_PROVIDER_FIELDS = [
"metar_read_zh",
"metar_read_en",
+ "taf_read_zh",
+ "taf_read_en",
+ "probability_read_zh",
+ "probability_read_en",
"predicted_max",
"range_low",
"range_high",
@@ -121,6 +129,10 @@ def _extract_json_string_field_fragment(raw_text: str, field: str) -> tuple[str,
_CITY_AI_TEXT_FIELDS = {
"metar_read_zh",
"metar_read_en",
+ "taf_read_zh",
+ "taf_read_en",
+ "probability_read_zh",
+ "probability_read_en",
"final_judgment_zh",
"final_judgment_en",
"reasoning_zh",
@@ -270,6 +282,10 @@ def _city_ai_response_example(unit: str) -> Dict[str, Any]:
return {
"metar_read_zh": f"最新观测显示 37.0{unit or '°C'},观测时间 04:30Z;风和云量暂未显示强降温信号,后续观测用于确认升温路径。",
"metar_read_en": f"The latest observation shows 37.0{unit or '°C'} at 04:30Z; wind and cloud signals do not yet show a strong cooling break, so later observations should confirm the warming path.",
+ "taf_read_zh": "TAF 预报 14-16Z BECMG 18012KT,午后风向切换为海风可能抑制升温。",
+ "taf_read_en": "TAF shows BECMG 18012KT during 14-16Z; afternoon onshore wind shift may cap warming.",
+ "probability_read_zh": f"概率分布偏右,最高概率桶 42-43{unit or '°C'}(~35%),上方尾部延至 45{unit or '°C'}。",
+ "probability_read_en": f"Distribution skews right; peak bucket 42-43{unit or '°C'} (~35%), upper tail to 45{unit or '°C'}.",
"final_judgment_zh": f"预计最高温暂以 43.0{unit or '°C'} 附近为中枢。",
"final_judgment_en": f"The expected daily high is centered near 43.0{unit or '°C'}.",
"predicted_max": 43.0,
@@ -290,6 +306,10 @@ def _city_ai_stream_response_example(unit: str) -> Dict[str, Any]:
return {
"metar_read_zh": f"最新 METAR 显示 37.0{unit or '°C'},报文时间 04:30Z;风和云量暂未显示强降温信号,后续报文用于确认升温路径。",
"metar_read_en": f"The latest METAR shows 37.0{unit or '°C'} at 04:30Z; wind and cloud signals do not yet show a strong cooling break, so later reports should confirm the warming path.",
+ "taf_read_zh": "TAF 预报 14-16Z BECMG 18012KT,午后风向转为海风可能抑制升温,需关注。",
+ "taf_read_en": "TAF shows BECMG 18012KT during 14-16Z; afternoon onshore wind shift may cap further warming.",
+ "probability_read_zh": f"概率分布偏右,最高概率桶 42-43{unit or '°C'}(~35%),上方尾部延至 45{unit or '°C'}。",
+ "probability_read_en": f"Distribution skews right; peak bucket 42-43{unit or '°C'} (~35%), upper tail extends to 45{unit or '°C'}.",
"predicted_max": 43.0,
"range_low": 42.0,
"range_high": 44.0,
diff --git a/web/scan_city_ai_prompt.py b/web/scan_city_ai_prompt.py
index 646a7e66..8ea17eaa 100644
--- a/web/scan_city_ai_prompt.py
+++ b/web/scan_city_ai_prompt.py
@@ -68,12 +68,14 @@ def build_city_ai_request_json(
"locale": normalized_locale,
"task": (
"Return strict JSON with: predicted_max, range_low, range_high, unit, confidence, "
- "final_judgment_zh, final_judgment_en, metar_read_zh, metar_read_en, "
+ "final_judgment_zh, final_judgment_en, metar_read_zh, metar_read_en, taf_read_zh, taf_read_en, probability_read_zh, probability_read_en, "
"reasoning_zh, reasoning_en, risks_zh, risks_en, model_cluster_note_zh, model_cluster_note_en. "
"Fill every *_zh field in Simplified Chinese and every *_en field in English in the same response. "
"Use this exact JSON object shape; do not return an array, markdown, or prose outside JSON. "
"Keep final_judgment one short decision sentence. metar_read must explain the latest observation source "
"with report/observation time, temperature, wind direction/speed, cloud/weather/visibility/dewpoint if available. "
+ "taf_read must interpret TAF for today's peak window impact (BECMG/TEMPO timing, wind shifts). If no TAF, write '无可用 TAF'/'No TAF available'. "
+ "probability_read must describe the probability distribution shape in 1 sentence (peak bucket, skew). "
"For wind, explicitly say whether the current wind tends to warm, cool, or be neutral for today's high, "
"and why in local city/station context. If mentioning cold/warm advection, name the wind direction or "
"direction shift responsible. If mentioning TAF risk, include the concrete TAF time window or say no "
@@ -203,12 +205,11 @@ def build_city_ai_stream_request(
context = _observation_prompt_context(ai_input)
is_airport_metar = context["is_airport_metar"]
role_label = "机场 METAR 解读与最高温预测员" if is_airport_metar else "官方观测站解读与最高温预测员"
- read_label = context["read_label_zh"]
source_instruction = context["instruction_zh"]
system_prompt = (
f"你是 PolyWeather 的{role_label}。"
"只返回一个紧凑 JSON object,不要 Markdown。"
- f"必须基于最新观测/报文独立判断该城市今日最高温,输出 metar_read_zh、metar_read_en、predicted_max、range_low、range_high、unit、confidence、final_judgment_zh、final_judgment_en、reasoning_zh、reasoning_en 字段,便于前端快速显示{read_label}和最高温预测;"
+ f"必须基于最新观测/报文独立判断该城市今日最高温,输出 metar_read_zh、metar_read_en、taf_read_zh、taf_read_en、probability_read_zh、probability_read_en、predicted_max、range_low、range_high、unit、confidence、final_judgment_zh、final_judgment_en、reasoning_zh、reasoning_en 字段;"
"模型一致性和风险清单由后端规则补齐,不要生成这些字段。"
f"预测方法:先看 model_cluster.sources 中各模型(含 DEB)的集中区间作为基线;"
"然后重点阅读 metar_context 和 current 中的观测信号——温度趋势、风向风速、湿度、云量、能见度——"
@@ -219,6 +220,8 @@ def build_city_ai_stream_request(
"如果 observation_anchor.is_airport_metar 为 false,不得使用 METAR、TAF、机场报文等称谓。"
"predicted_max 是你的独立预测值(float),range_low/range_high 是预测区间,unit 为温度单位,confidence 为 low/medium/high。"
"final_judgment 用一句话给出今日最高温结论。reasoning 必须解释你相对于模型集群基线做了哪种修正及原因。"
+ "taf_read_zh/en: 如果 city_snapshot.taf 有有效内容,用 1-2 句解读机场预报中对今日峰值窗口有影响的变化(BECMG/TEMPO 时间窗、风向切换、云量变化);如果无 TAF 或 TAF 不含今日白天有效时段,写「无可用 TAF」/「No TAF available」。"
+ "probability_read_zh/en: 如果 city_snapshot.probability 有分布数据,用 1 句描述概率分布形态——最高概率桶落在哪、分布偏左/偏右/对称。"
"所有 *_zh 字段写简体中文,所有 *_en 字段写英文,不得留空。"
"不要写交易建议、BUY/SELL、Kelly 或套利。"
)
@@ -236,10 +239,12 @@ def build_city_ai_stream_request(
{
"locale": normalized_locale,
"task": (
- "Return JSON keys in this exact order: metar_read_zh, metar_read_en, predicted_max, range_low, range_high, unit, confidence, final_judgment_zh, final_judgment_en, reasoning_zh, reasoning_en. "
+ "Return JSON keys in this exact order: metar_read_zh, metar_read_en, taf_read_zh, taf_read_en, probability_read_zh, probability_read_en, predicted_max, range_low, range_high, unit, confidence, final_judgment_zh, final_judgment_en, reasoning_zh, reasoning_en. "
"predicted_max must be your independent float prediction, based on model cluster baseline adjusted by the latest METAR/observation signals. "
"Do not copy DEB directly \u2014 use the full model spread + your own reading of wind, cloud, temperature trend from the bulletin. "
"reasoning must explain what adjustment you made relative to the model cluster and why. "
+ "taf_read must interpret TAF for peak window impact if available. "
+ "probability_read must describe the probability distribution shape in 1 sentence. "
"Do not return risks or model_cluster_note. Keep it compact. "
"Return exactly one JSON object and no markdown."
),
diff --git a/web/scan_terminal_ai_compact.py b/web/scan_terminal_ai_compact.py
index 241ad87c..90c63d1a 100644
--- a/web/scan_terminal_ai_compact.py
+++ b/web/scan_terminal_ai_compact.py
@@ -421,3 +421,46 @@ def build_scan_ai_prompt(payload: Dict[str, Any], *, max_rows: int) -> Dict[str,
"sent_contracts": sent_contracts,
},
}
+
+
+def _compact_probability_context(probabilities: Any, deb: Any, unit: str) -> dict:
+ if not isinstance(probabilities, dict):
+ return {}
+ dist = probabilities.get("distribution")
+ if not isinstance(dist, list) or not dist:
+ return {}
+ top_buckets = sorted(
+ [b for b in dist if isinstance(b, dict) and b.get("probability")],
+ key=lambda b: float(b.get("probability", 0)),
+ reverse=True,
+ )[:3]
+ compact = {
+ "top_buckets": [
+ {
+ "label": b.get("label", ""),
+ "prob": round(float(b.get("probability", 0)) * 100),
+ }
+ for b in top_buckets
+ ],
+ }
+ mu = probabilities.get("mu")
+ if mu is not None:
+ compact["mu"] = mu
+ spread = (
+ round(float(probabilities.get("calibrated_sigma") or 0), 1)
+ or round(float(probabilities.get("raw_sigma") or 0), 1)
+ or None
+ )
+ if spread is not None:
+ compact["sigma"] = spread
+ deb_val = deb.get("prediction") if isinstance(deb, dict) else None
+ if deb_val is not None and mu is not None:
+ if deb_val > mu:
+ compact["skew"] = "right"
+ elif deb_val < mu:
+ compact["skew"] = "left"
+ else:
+ compact["skew"] = "centered"
+ if unit:
+ compact["unit"] = unit
+ return compact
diff --git a/web/scan_terminal_service.py b/web/scan_terminal_service.py
index e527fe9e..d02b8895 100644
--- a/web/scan_terminal_service.py
+++ b/web/scan_terminal_service.py
@@ -51,6 +51,7 @@ from web.scan_terminal_ai_compact import (
_compact_hourly_context,
_compact_intraday_context,
_compact_observation_points,
+ _compact_probability_context,
_compact_taf_context,
_compact_vertical_context,
build_scan_ai_prompt,
@@ -177,7 +178,7 @@ SCAN_CITY_AI_RETRY_ON_STREAM_PARSE_ERROR = str(
).strip().lower() in {"1", "true", "yes", "on"}
SCAN_AI_CACHE_TTL_SEC = max(
30,
- int(os.getenv("POLYWEATHER_SCAN_AI_CACHE_TTL_SEC", "1800")),
+ int(os.getenv("POLYWEATHER_SCAN_AI_CACHE_TTL_SEC", "3600")),
)
SCAN_AI_MAX_ROWS = _env_int("POLYWEATHER_SCAN_AI_MAX_ROWS", 40, min_value=1)
SCAN_AI_MAX_TOKENS = _env_int(
@@ -194,7 +195,7 @@ SCAN_CITY_AI_MAX_TOKENS = _env_int(
)
SCAN_CITY_AI_STREAM_MAX_TOKENS = _env_int(
"POLYWEATHER_SCAN_CITY_AI_STREAM_MAX_TOKENS",
- min(SCAN_CITY_AI_MAX_TOKENS, 900),
+ min(SCAN_CITY_AI_MAX_TOKENS, 1200),
min_value=400,
max_value=64000,
)
@@ -393,6 +394,11 @@ def _build_city_ai_prompt(data: Dict[str, Any]) -> Dict[str, Any]:
"station_label": airport_current.get("station_label"),
},
"taf": _compact_taf_context(data.get("taf")),
+ "probability": _compact_probability_context(
+ data.get("probabilities") if isinstance(data.get("probabilities"), dict) else None,
+ data.get("deb") if isinstance(data.get("deb"), dict) else None,
+ data.get("temp_symbol"),
+ ),
"vertical_profile_signal": _compact_vertical_context(
data.get("vertical_profile_signal")
),
@@ -425,47 +431,39 @@ def _scan_city_ai_cache_key(ai_input: Dict[str, Any]) -> str:
observation_anchor = ai_input.get("observation_anchor") if isinstance(ai_input.get("observation_anchor"), dict) else {}
is_airport_metar = observation_anchor.get("is_airport_metar") is not False
airport_current = ai_input.get("airport_current") if isinstance(ai_input.get("airport_current"), dict) else {}
- current_obs = ai_input.get("current") if isinstance(ai_input.get("current"), dict) else {}
metar_context = ai_input.get("metar_context") if isinstance(ai_input.get("metar_context"), dict) else {}
- observation_obs = (
- ai_input.get("metar_today_obs") or ai_input.get("metar_recent_obs") or []
- if is_airport_metar
- else ai_input.get("settlement_today_obs") or ai_input.get("settlement_recent_obs") or []
- )
- observation_fingerprint = {
- "stale_for_today": metar_context.get("stale_for_today"),
- "last_observation_time": metar_context.get("last_observation_time"),
- "last_time": metar_context.get("last_time"),
- "last_temp": metar_context.get("last_temp"),
- "max_time": metar_context.get("max_time"),
- "max_temp": metar_context.get("max_temp"),
- "airport_obs_time": airport_current.get("obs_time"),
- "airport_report_time": airport_current.get("report_time"),
- "airport_receipt_time": airport_current.get("receipt_time"),
- "airport_temp": airport_current.get("temp"),
- "airport_max_so_far": airport_current.get("max_so_far"),
- "current_obs_time": current_obs.get("obs_time"),
- "current_report_time": current_obs.get("report_time"),
- "current_temp": current_obs.get("temp"),
- "current_max_so_far": current_obs.get("max_so_far"),
- }
key_payload = {
"prompt_version": SCAN_CITY_AI_PROMPT_VERSION,
- "schema_version": ai_input.get("schema_version"),
- "model": SCAN_CITY_AI_MODEL,
"city": ai_input.get("city"),
"local_date": ai_input.get("local_date"),
- "deb": (ai_input.get("deb") or {}).get("prediction") if isinstance(ai_input.get("deb"), dict) else None,
- "observation_source": observation_anchor.get("source") or ("metar" if is_airport_metar else "official"),
"station": observation_anchor.get("station_code"),
- "metar": airport_current.get("raw_metar") if is_airport_metar else None,
- "observation_fingerprint": observation_fingerprint,
- "obs": observation_obs,
+ "raw_metar": airport_current.get("raw_metar") if is_airport_metar else None,
+ "obs_time": airport_current.get("obs_time") or metar_context.get("last_observation_time"),
+ "stale_for_today": metar_context.get("stale_for_today"),
}
raw = json.dumps(key_payload, sort_keys=True, ensure_ascii=False, default=str)
return "city-ai:" + hashlib.sha256(raw.encode("utf-8")).hexdigest()
+def _quick_metar_cache_key(data: Dict[str, Any]) -> str:
+ airport_current = data.get("airport_current") if isinstance(data.get("airport_current"), dict) else {}
+ current = data.get("current") if isinstance(data.get("current"), dict) else {}
+ observation_anchor = data.get("observation_anchor") if isinstance(data.get("observation_anchor"), dict) else {}
+ raw_metar = airport_current.get("raw_metar") or current.get("raw_metar")
+ obs_time = airport_current.get("obs_time") or current.get("obs_time")
+ if raw_metar and obs_time:
+ finger = {
+ "city": data.get("name"),
+ "raw_metar": raw_metar,
+ "obs_time": obs_time,
+ "station": observation_anchor.get("station_code"),
+ "prompt_version": SCAN_CITY_AI_PROMPT_VERSION,
+ }
+ raw = json.dumps(finger, sort_keys=True, ensure_ascii=False, default=str)
+ return "city-ai:" + hashlib.sha256(raw.encode("utf-8")).hexdigest()
+ return ""
+
+
def _sse_event(event: str, payload: Dict[str, Any]) -> str:
return (
f"event: {event}\n"
@@ -492,15 +490,19 @@ def _cache_city_ai_payload(
data: Dict[str, Any],
generated_at: str,
ai_raw: Dict[str, Any],
+ quick_key: str = "",
) -> None:
+ entry = {
+ "expires_at": time.time() + SCAN_AI_CACHE_TTL_SEC,
+ "generated_at": generated_at,
+ "city": data.get("name"),
+ "city_display_name": data.get("display_name"),
+ "payload": ai_raw,
+ }
with _SCAN_CITY_AI_CACHE_LOCK:
- _SCAN_CITY_AI_CACHE[cache_key] = {
- "expires_at": time.time() + SCAN_AI_CACHE_TTL_SEC,
- "generated_at": generated_at,
- "city": data.get("name"),
- "city_display_name": data.get("display_name"),
- "payload": ai_raw,
- }
+ _SCAN_CITY_AI_CACHE[cache_key] = entry
+ if quick_key and quick_key != cache_key:
+ _SCAN_CITY_AI_CACHE[quick_key] = entry
def _is_city_ai_fallback(ai_raw: Any) -> bool:
@@ -590,9 +592,12 @@ def stream_scan_city_ai_forecast_payload(
)
ai_input = _build_city_ai_prompt(data)
cache_key = _scan_city_ai_cache_key(ai_input)
+ quick_key = _quick_metar_cache_key(data)
if not force_refresh:
with _SCAN_CITY_AI_CACHE_LOCK:
- cached = _SCAN_CITY_AI_CACHE.get(cache_key)
+ cached = _SCAN_CITY_AI_CACHE.get(cache_key) or (
+ _SCAN_CITY_AI_CACHE.get(quick_key) if quick_key else None
+ )
if cached and cached.get("expires_at", 0) >= time.time():
yield _sse_event(
"final",
@@ -803,6 +808,7 @@ def stream_scan_city_ai_forecast_payload(
data=data,
generated_at=generated_at,
ai_raw=ai_raw,
+ quick_key=quick_key,
)
yield _sse_event(
"final",
@@ -914,9 +920,12 @@ def build_scan_city_ai_forecast_payload(
)
ai_input = _build_city_ai_prompt(data)
cache_key = _scan_city_ai_cache_key(ai_input)
+ quick_key = _quick_metar_cache_key(data)
if not force_refresh:
with _SCAN_CITY_AI_CACHE_LOCK:
- cached = _SCAN_CITY_AI_CACHE.get(cache_key)
+ cached = _SCAN_CITY_AI_CACHE.get(cache_key) or (
+ _SCAN_CITY_AI_CACHE.get(quick_key) if quick_key else None
+ )
if cached and cached.get("expires_at", 0) >= time.time():
logger.info(
"scan city AI forecast cache hit city={} model={}",
@@ -1058,6 +1067,7 @@ def build_scan_city_ai_forecast_payload(
data=data,
generated_at=generated_at,
ai_raw=ai_raw,
+ quick_key=quick_key,
)
logger.info(
"scan city AI forecast complete city={} duration_ms={} model={} confidence={}",
diff --git a/web/services/anomaly_detection.py b/web/services/anomaly_detection.py
new file mode 100644
index 00000000..c96bd773
--- /dev/null
+++ b/web/services/anomaly_detection.py
@@ -0,0 +1,102 @@
+"""Anomaly detection — pure math, no AI call.
+
+Flags cities where current observations deviate from model predictions.
+"""
+
+from __future__ import annotations
+
+from typing import Any, Dict, List, Optional
+
+from web.scan_city_ai_helpers import _safe_float
+
+
+def _check_city_anomaly(
+ data: Dict[str, Any],
+ *,
+ high_temp_threshold: float = 2.0,
+) -> Optional[Dict[str, Any]]:
+ """Return anomaly flag if current observation breaks model cluster bounds."""
+ current = data.get("current") if isinstance(data.get("current"), dict) else {}
+ airport = data.get("airport_current") if isinstance(data.get("airport_current"), dict) else {}
+ multi = data.get("multi_model") if isinstance(data.get("multi_model"), dict) else {}
+ deb = data.get("deb") if isinstance(data.get("deb"), dict) else {}
+
+ observed = _safe_float(current.get("temp") or airport.get("temp"))
+ if observed is None:
+ return None
+
+ model_highs = [
+ _safe_float(v)
+ for v in multi.values()
+ if _safe_float(v) is not None
+ ]
+ deb_pred = _safe_float(deb.get("prediction"))
+ if deb_pred is not None:
+ model_highs.append(deb_pred)
+
+ if not model_highs:
+ return None
+
+ model_max = max(model_highs)
+ model_min = min(model_highs)
+ model_median = sorted(model_highs)[len(model_highs) // 2]
+
+ delta_above_max = observed - model_max
+ delta_below_min = model_min - observed
+ delta_from_median = observed - model_median
+
+ anomaly: Optional[Dict[str, Any]] = None
+
+ if delta_above_max > high_temp_threshold:
+ anomaly = {
+ "level": "breakout_above",
+ "observed": observed,
+ "model_max": model_max,
+ "delta": round(delta_above_max, 1),
+ "model_count": len(model_highs),
+ }
+ elif delta_below_min > high_temp_threshold:
+ anomaly = {
+ "level": "breakout_below",
+ "observed": observed,
+ "model_min": model_min,
+ "delta": round(delta_below_min, 1),
+ "model_count": len(model_highs),
+ }
+ elif abs(delta_from_median) > 1.5:
+ anomaly = {
+ "level": "deviation",
+ "observed": observed,
+ "model_median": model_median,
+ "delta": round(delta_from_median, 1),
+ "model_count": len(model_highs),
+ }
+
+ if anomaly:
+ anomaly.update(
+ {
+ "city": data.get("name") or data.get("city"),
+ "local_date": data.get("local_date"),
+ "temp_unit": data.get("temp_symbol", "°C"),
+ "deb_prediction": deb_pred,
+ }
+ )
+ return anomaly
+
+
+def detect_scan_terminal_anomalies(
+ rows: List[Dict[str, Any]],
+ *,
+ high_temp_threshold: float = 2.0,
+) -> List[Dict[str, Any]]:
+ """Scan all terminal rows and return anomaly flags."""
+ anomalies = []
+ for row in rows:
+ if not isinstance(row, dict):
+ continue
+ city_data = row.get("city_data") or row
+ flag = _check_city_anomaly(city_data, high_temp_threshold=high_temp_threshold)
+ if flag:
+ flag["row_id"] = row.get("row_id") or row.get("id")
+ anomalies.append(flag)
+ return anomalies
diff --git a/web/services/market_overview_api.py b/web/services/market_overview_api.py
new file mode 100644
index 00000000..9b14f2ba
--- /dev/null
+++ b/web/services/market_overview_api.py
@@ -0,0 +1,196 @@
+"""Market overview — AI summary of all scan terminal rows, cached 10 min."""
+
+from __future__ import annotations
+
+import hashlib
+import json
+import threading
+import time
+from datetime import datetime
+from typing import Any, Dict, List
+
+from loguru import logger
+
+from web.scan_city_ai_helpers import _safe_float
+from web.scan_terminal_service import (
+ SCAN_AI_BASE_URL,
+ SCAN_CITY_AI_MODEL,
+ SCAN_CITY_AI_TIMEOUT_SEC,
+ _scan_ai_api_key,
+)
+
+_OVERVIEW_CACHE: Dict[str, Dict[str, Any]] = {}
+_OVERVIEW_CACHE_LOCK = threading.Lock()
+_OVERVIEW_MAX_TOKENS = 600
+_OVERVIEW_CACHE_TTL_SEC = 600
+
+
+def _build_overview_ai_request(
+ rows: List[Dict[str, Any]],
+ locale: str,
+) -> Dict[str, Any]:
+ cities = []
+ for row in rows:
+ if not isinstance(row, dict):
+ continue
+ city = row.get("city") or row.get("name") or ""
+ if not city:
+ continue
+ model_cluster = row.get("model_cluster") if isinstance(row.get("model_cluster"), dict) else {}
+ sources = model_cluster.get("sources") if isinstance(model_cluster.get("sources"), list) else []
+ values = [
+ _safe_float(s.get("value"))
+ for s in sources
+ if isinstance(s, dict) and _safe_float(s.get("value")) is not None
+ ]
+ deb_val = _safe_float(row.get("deb_prediction") or (row.get("deb") or {}).get("prediction"))
+ cities.append(
+ {
+ "city": str(city),
+ "display_name": row.get("display_name") or str(city),
+ "local_date": row.get("local_date", ""),
+ "deb": deb_val,
+ "model_min": min(values) if values else None,
+ "model_max": max(values) if values else None,
+ "model_count": len(values),
+ "current_temp": _safe_float(row.get("current_temp") or (row.get("current") or {}).get("temp")),
+ "max_so_far": _safe_float(row.get("current_max_so_far") or row.get("max_so_far") or (row.get("current") or {}).get("max_so_far")),
+ "risk_level": row.get("risk_level", ""),
+ "temp_unit": row.get("temp_unit") or row.get("temp_symbol") or "°C",
+ }
+ )
+
+ system_prompt = (
+ "你是 PolyWeather 的天气市场概览员。基于全部城市的扫描数据,写一段今日市场概览。"
+ "用 3-5 句概括:整体模型一致性、最值得关注的城市(模型分歧大或实测偏离集群)、异常信号。"
+ "highlights 最多 5 个城市,每个城市一句话点出关键信号。"
+ "只返回 JSON object,不要 Markdown。所有 *_zh 字段写简体中文,*_en 字段写英文。"
+ )
+ task = (
+ "Return JSON: overview_zh, overview_en, highlights (array of {city, note_zh, note_en}, max 5). "
+ "overview: 3-5 sentences covering model consensus, top divergence cities, anomalies. "
+ "highlights: per-city one-sentence signal. Keep compact."
+ )
+
+ return {
+ "model": SCAN_CITY_AI_MODEL,
+ "temperature": 0.3,
+ "max_tokens": _OVERVIEW_MAX_TOKENS,
+ "response_format": {"type": "json_object"},
+ "messages": [
+ {"role": "system", "content": system_prompt},
+ {
+ "role": "user",
+ "content": json.dumps(
+ {
+ "locale": locale,
+ "task": task,
+ "city_count": len(cities),
+ "cities": cities,
+ },
+ ensure_ascii=False,
+ default=str,
+ ),
+ },
+ ],
+ }
+
+
+def _cache_key(rows: List[Dict[str, Any]], locale: str) -> str:
+ finger = {
+ "city_ids": sorted(
+ row.get("city") or row.get("name") or ""
+ for row in rows
+ if isinstance(row, dict)
+ ),
+ "locale": locale,
+ }
+ raw = json.dumps(finger, sort_keys=True, ensure_ascii=False, default=str)
+ return "overview:" + hashlib.sha256(raw.encode("utf-8")).hexdigest()
+
+
+def build_market_overview_payload(
+ rows: List[Dict[str, Any]],
+ *,
+ locale: str = "zh-CN",
+ force_refresh: bool = False,
+) -> Dict[str, Any]:
+ if not rows:
+ return {"overview_zh": "", "overview_en": "", "highlights": [], "generated_at": None}
+
+ key = _cache_key(rows, locale)
+ if not force_refresh:
+ with _OVERVIEW_CACHE_LOCK:
+ cached = _OVERVIEW_CACHE.get(key)
+ if cached and cached.get("expires_at", 0) >= time.time():
+ return cached["payload"]
+
+ api_key = _scan_ai_api_key()
+ if not api_key:
+ return {
+ "overview_zh": "AI 概览不可用(未配置 API Key)",
+ "overview_en": "AI overview unavailable (API key not configured)",
+ "highlights": [],
+ "generated_at": datetime.utcnow().isoformat() + "Z",
+ }
+
+ import httpx
+
+ request_json = _build_overview_ai_request(rows, locale)
+ generated_at = datetime.utcnow().isoformat() + "Z"
+ started = time.perf_counter()
+
+ try:
+ response = httpx.post(
+ f"{SCAN_AI_BASE_URL}/chat/completions",
+ json=request_json,
+ headers={
+ "Authorization": f"Bearer {api_key}",
+ "Content-Type": "application/json",
+ },
+ timeout=min(SCAN_CITY_AI_TIMEOUT_SEC, 15),
+ )
+ response.raise_for_status()
+ result = response.json()
+ content = ((result.get("choices") or [{}])[0].get("message") or {}).get("content") or "{}"
+ parsed = json.loads(content) if isinstance(content, str) else content
+ if not isinstance(parsed, dict):
+ raise ValueError("AI returned non-dict overview")
+
+ payload: Dict[str, Any] = {
+ "overview_zh": str(parsed.get("overview_zh") or parsed.get("overview_en") or ""),
+ "overview_en": str(parsed.get("overview_en") or parsed.get("overview_zh") or ""),
+ "highlights": [
+ {
+ "city": str(h.get("city", "")),
+ "note_zh": str(h.get("note_zh", "")),
+ "note_en": str(h.get("note_en", "")),
+ }
+ for h in (parsed.get("highlights") if isinstance(parsed.get("highlights"), list) else [])
+ if isinstance(h, dict)
+ ][:5],
+ "generated_at": generated_at,
+ }
+ except Exception as exc:
+ logger.warning("Market overview AI failed: {}", exc)
+ payload = {
+ "overview_zh": "市场概览暂时无法生成,请稍后刷新。",
+ "overview_en": "Market overview temporarily unavailable, please refresh later.",
+ "highlights": [],
+ "generated_at": generated_at,
+ }
+
+ duration_ms = int((time.perf_counter() - started) * 1000)
+ logger.info(
+ "market_overview cities={} locale={} duration_ms={} cached={}",
+ len(rows),
+ locale,
+ duration_ms,
+ False,
+ )
+
+ entry = {"expires_at": time.time() + _OVERVIEW_CACHE_TTL_SEC, "payload": payload}
+ with _OVERVIEW_CACHE_LOCK:
+ _OVERVIEW_CACHE[key] = entry
+
+ return payload
diff --git a/web/services/scan_api.py b/web/services/scan_api.py
index 393df4a8..7143f97c 100644
--- a/web/services/scan_api.py
+++ b/web/services/scan_api.py
@@ -8,6 +8,7 @@ from fastapi import HTTPException, Request
from fastapi.concurrency import run_in_threadpool
from fastapi.responses import StreamingResponse
+from web.services.market_overview_api import build_market_overview_payload
import web.routes as legacy_routes
@@ -109,3 +110,21 @@ async def get_scan_city_ai_stream_response(request: Request) -> StreamingRespons
"X-Accel-Buffering": "no",
},
)
+
+
+async def get_scan_terminal_overview_payload(request: Request) -> Dict[str, Any]:
+ try:
+ body = await request.json()
+ except Exception:
+ body = {}
+ if not isinstance(body, dict):
+ raise HTTPException(status_code=400, detail="Invalid JSON body")
+ rows = body.get("rows") if isinstance(body.get("rows"), list) else []
+ locale = str(body.get("locale") or "zh-CN").strip()
+ force_refresh = str(body.get("force_refresh") or "false").strip().lower() in {"1", "true", "yes"}
+ return await run_in_threadpool(
+ build_market_overview_payload,
+ rows,
+ locale=locale,
+ force_refresh=force_refresh,
+ )