From 2ee00f8016789e7b467d911f480efd667f169fe1 Mon Sep 17 00:00:00 2001 From: "2569718930@qq.com" <2569718930@qq.com> Date: Thu, 14 May 2026 22:36:19 +0800 Subject: [PATCH] =?UTF-8?q?MiMo=20AI=20=E8=83=BD=E5=8A=9B=E6=89=A9?= =?UTF-8?q?=E5=B1=95=EF=BC=9ATAF=E8=A7=A3=E8=AF=BB=E3=80=81=E6=A6=82?= =?UTF-8?q?=E7=8E=87=E5=88=86=E5=B8=83=E8=A7=A3=E8=AF=BB=E3=80=81=E5=BC=82?= =?UTF-8?q?=E5=B8=B8=E6=A3=80=E6=B5=8B=E3=80=81=E5=B8=82=E5=9C=BA=E6=A6=82?= =?UTF-8?q?=E8=A7=88?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit AI 解读字段扩展: - 新增 taf_read_zh/en:解读机场预报中影响今日峰值窗口的变化 - 新增 probability_read_zh/en:描述概率分布形态(最高桶、偏左/偏右) - stream max_tokens 900→1200 容纳新输出字段 - 缓存 key 简化为 METAR原文+观测时间,大幅提升命中率 - 兜底函数补全 TAF 和概率字段的确定性生成 异常检测: - 纯数学计算,零 AI 延迟:实测温度 vs 全部模型预测上下限 - 三级告警:breakout_above / breakout_below / deviation 市场概览: - 新增 POST /api/scan/terminal/overview(MiMo 批量解读,缓存10分钟) - 前端 MarketOverviewBanner 可折叠横幅(顶栏与标签栏之间) - 移动端适配 640px/768px 断点,暗色/亮色双主题 Scope-risk: MEDIUM — 170 测试通过,TypeScript 零错误,ruff 零告警 Tested: python -m pytest -q (170 passed), npx tsc --noEmit (0 errors), ruff check . --- .../dashboard/ScanTerminalDashboard.tsx | 14 ++ .../MarketOverviewBanner.module.css | 192 +++++++++++++++++ .../scan-terminal/MarketOverviewBanner.tsx | 132 ++++++++++++ tests/test_web_observability.py | 17 +- web/routers/scan.py | 6 + web/scan_city_ai_fallback.py | 32 +++ web/scan_city_ai_helpers.py | 20 ++ web/scan_city_ai_prompt.py | 13 +- web/scan_terminal_ai_compact.py | 43 ++++ web/scan_terminal_service.py | 92 ++++---- web/services/anomaly_detection.py | 102 +++++++++ web/services/market_overview_api.py | 196 ++++++++++++++++++ web/services/scan_api.py | 19 ++ 13 files changed, 823 insertions(+), 55 deletions(-) create mode 100644 frontend/components/dashboard/scan-terminal/MarketOverviewBanner.module.css create mode 100644 frontend/components/dashboard/scan-terminal/MarketOverviewBanner.tsx create mode 100644 web/services/anomaly_detection.py create mode 100644 web/services/market_overview_api.py diff --git a/frontend/components/dashboard/ScanTerminalDashboard.tsx b/frontend/components/dashboard/ScanTerminalDashboard.tsx index 6c01cacf..8ec1a8c8 100644 --- a/frontend/components/dashboard/ScanTerminalDashboard.tsx +++ b/frontend/components/dashboard/ScanTerminalDashboard.tsx @@ -46,6 +46,14 @@ import { useUserLocalClock, } from "@/components/dashboard/scan-terminal/use-scan-terminal-ui-state"; import { useRelativeTime } from "@/hooks/useRelativeTime"; +const MarketOverviewBanner = dynamic( + () => + import( + "@/components/dashboard/scan-terminal/MarketOverviewBanner" + ).then((m) => m.MarketOverviewBanner), + { ssr: false }, +); + const MonitorPanel = dynamic( () => import("@/components/dashboard/monitoring/MonitorPanel"), { ssr: false }, @@ -423,6 +431,12 @@ function ScanTerminalScreen() { /> ) : null} + +
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, + )