MiMo AI 能力扩展:TAF解读、概率分布解读、异常检测、市场概览

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 .
This commit is contained in:
2569718930@qq.com
2026-05-14 22:36:19 +08:00
parent 6c08a68413
commit 2ee00f8016
13 changed files with 823 additions and 55 deletions
@@ -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}
<MarketOverviewBanner
isEn={isEn}
isPro={isPro}
rows={timeSortedRows}
/>
<section className="scan-list-section">
<div className="scan-list-header">
<div className="scan-list-tabs" role="tablist" aria-label={isEn ? "Content view" : "内容视图"}>
@@ -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);
}
@@ -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<OverviewPayload | null>(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 (
<div className={clsx(styles.root, styles.loading)}>
<Sparkles size={14} className={styles.icon} />
<span>{isEn ? "AI is generating market overview…" : "AI 正在生成市场概览…"}</span>
</div>
);
}
if (error && !data) return null;
const overviewText = data ? (isEn ? data.overview_en : data.overview_zh) : "";
const highlights = data?.highlights ?? [];
return (
<div className={clsx(styles.root, collapsed && styles.collapsed)}>
<button
type="button"
className={styles.header}
onClick={() => setCollapsed((c) => !c)}
aria-label={isEn ? "Toggle market overview" : "切换市场概览"}
>
<span className={styles.badge}>
<Sparkles size={13} />
{isEn ? "AI Overview" : "AI 概览"}
</span>
<span className={styles.preview}>
{collapsed && overviewText ? overviewText.slice(0, isEn ? 120 : 60) + "…" : ""}
</span>
<span className={styles.toggle}>
{collapsed ? <ChevronDown size={16} /> : <ChevronUp size={16} />}
</span>
</button>
{!collapsed && (
<div className={styles.body}>
<p className={styles.summary}>{overviewText}</p>
{highlights.length > 0 && (
<ul className={styles.highlights}>
{highlights.map((h) => (
<li key={h.city} className={styles.highlightItem}>
<strong>{h.city}</strong>
<span>{isEn ? h.note_en : h.note_zh}</span>
</li>
))}
</ul>
)}
{data?.generated_at && (
<time className={styles.time} dateTime={data.generated_at}>
{isEn ? "Generated " : "生成于 "}
{new Date(data.generated_at).toLocaleTimeString(isEn ? "en-US" : "zh-CN", {
hour: "2-digit",
minute: "2-digit",
})}
</time>
)}
</div>
)}
</div>
);
}
+7 -10
View File
@@ -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
+6
View File
@@ -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)
+32
View File
@@ -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,
+20
View File
@@ -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,
+9 -4
View File
@@ -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."
),
+43
View File
@@ -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
+51 -41
View File
@@ -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={}",
+102
View File
@@ -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
+196
View File
@@ -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
+19
View File
@@ -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,
)