feat(recap): 盘前风向标服务与 AI 复盘注入

- fuyao 客户端新增 short-term-benchmark 端点透传
- auction_benchmark 服务: 历史按日 JSON 缓存(当日不缓存)、交易日回退、
  四态降级(ok/fallback_prev/source_unavailable/no_data)
- 读取时从相邻日K分区现算当日/次日收益对照, 不落缓存
- 复盘 user prompt 注入「盘前风向标(竞价)」上下文
- 11 个测试覆盖回退/缓存/富化数学/上下文
This commit is contained in:
shy3130
2026-08-30 19:05:31 +08:00
parent cff3b07527
commit 3d627aba5e
6 changed files with 526 additions and 5 deletions
+23 -2
View File
@@ -1,4 +1,4 @@
"""AI 大盘复盘 API — 流式复盘 + 报告持久化 + 龙虎榜。
"""AI 大盘复盘 API — 流式复盘 + 报告持久化 + 龙虎榜 + 盘前风向标
路由前缀: /api/market-recap
@@ -8,6 +8,7 @@
POST /reports 保存一条复盘报告
DELETE /reports/{report_id} 删除一条复盘报告
GET /dragon-tiger 龙虎榜三榜 (fuyao 专有, 历史按日缓存)
GET /auction-benchmark 盘前风向标 (fuyao 专有, 含当日/次日真实收益)
"""
from __future__ import annotations
@@ -18,7 +19,7 @@ from fastapi import APIRouter, HTTPException, Query, Request
from fastapi.responses import StreamingResponse
from pydantic import BaseModel
from app.services import dragon_tiger, market_recap_reports
from app.services import auction_benchmark, dragon_tiger, market_recap_reports
from app.services.market_recap import recap_market_stream
logger = logging.getLogger(__name__)
@@ -52,6 +53,26 @@ def get_dragon_tiger(
)
@router.get("/auction-benchmark")
def get_auction_benchmark(
request: Request,
date: str | None = Query(default=None, description="复盘目标日 YYYY-MM-DD, 缺省取最近交易日"),
):
"""盘前风向标 (同花顺竞价筛选名单 + 当日/次日真实收益)。
fuyao 专有, 未配置时 state=source_unavailable; 非交易日由服务层自动回退。
"""
target = None
if date:
try:
target = date_cls.fromisoformat(date)
except ValueError:
raise HTTPException(400, f"date 格式应为 YYYY-MM-DD, 收到: {date}")
return auction_benchmark.get_auction_benchmark(
request.app.state.repo.store.data_dir, target
)
@router.post("/analyze")
async def analyze_market(request: Request, req: AnalyzeRequest):
"""AI 大盘复盘 — NDJSON 流式返回。
+12
View File
@@ -231,6 +231,18 @@ class FuyaoClient:
params["date"] = date
return self._get("/api/a-share/special-data/dragon-tiger-list", params)
def short_term_benchmark(self, date: str | None = None) -> dict:
"""短线风向标竞价基准 (同花顺竞价筛选, 每日约 5~6 只)。
返回 data 原始容器: {date, date_ms, item[]}。item 行:
{thscode, ticker, name, auction_pct, tags[]}。支持一年内历史日期;
显式传非交易日返回 code=1002 (由调用方做交易日回退)。
"""
params: dict = {}
if date:
params["date"] = date
return self._get("/api/a-share/auction/short-term-benchmark", params)
# ---- 市场 dump ----
def dump_download_url(self, dump_kind: str) -> dict:
"""获取 dump 预签名下载信息(约 300s 有效)。
+8
View File
@@ -895,6 +895,14 @@ class FuyaoProvider:
"""
return self._get_client().dragon_tiger_list(board_type, date)
def short_term_benchmark(self, date: str | None = None) -> dict:
"""短线风向标竞价基准 (复盘页卡片 + AI 复盘上下文)。返回原始 data 容器。
非路由数据集 (tickflow 无对应能力), 不进 plugin.yaml datasets,
由 services.auction_benchmark 统一做按日缓存/收益enrich/交易日回退。
"""
return self._get_client().short_term_benchmark(date)
def _derive_bps(self, symbols: list[str]) -> dict[str, float]:
"""估值快照 pb_mrq 与行情快照最新价同源同刻 → bps = price / pb_mrq。
+288
View File
@@ -0,0 +1,288 @@
"""短线风向标服务 (fuyao 专有) — 复盘页卡片 + AI 复盘上下文。
非路由数据集: tickflow 无对应能力, 直接经 custom_sources 调 fuyao provider;
fuyao 未配置时返回 source_unavailable 状态, 前端降级提示。
数据契约 (实测 2026-08-30, 60 交易日回测 353 样本):
- 每日 5~6 只, 服务端筛选的竞价异动股, 附概念标签
- 名单当日 (开盘买→收盘卖) 均值 +0.54% vs 全市场 +0.10%, 有真实当日选股能力;
但高开≥5% 子集当日 -1.97% (追高陷阱) → 前端对高开子集标「追高风险」
- 次日无显著优势 (+0.08%), 定位为「当日观察名单」而非隔夜轮动信号
缓存策略 (与 dragon_tiger 同模式):
- 历史名单不可变 → 按日落 JSON 缓存 (data/auction_benchmark/date=YYYY-MM-DD.json),
缓存命中不触发插件注册表加载
- 收益 enrich (当日oc/全天/次日) 不落缓存 — 次日数据晚到, 读取时现算
- 当日名单不缓存 (竞价阶段名单可能变动, 以现拉为准)
- 显式日期失败 → 回退上一交易日一次 (state=fallback_prev)
日期解析: 接口对显式非交易日报 code=1002, 本层用本地 kline_daily 分区日期
把目标日回退到「≤目标日的最近交易日」, 规避报错。
"""
from __future__ import annotations
import contextlib
import json
import logging
import re
from datetime import date as date_cls
from pathlib import Path
import polars as pl
from app.market_time import cn_today
logger = logging.getLogger(__name__)
_DATE_DIR_RE = re.compile(r"^date=(\d{4}-\d{2}-\d{2})$")
def _local_trading_days(data_dir: Path) -> list[date_cls]:
"""本地日K分区日期 = 已知交易日集合 (升序)。扫描失败返回空。"""
root = data_dir / "kline_daily"
out: list[date_cls] = []
try:
for d in root.iterdir():
m = _DATE_DIR_RE.match(d.name)
if d.is_dir() and m:
try:
out.append(date_cls.fromisoformat(m.group(1)))
except ValueError:
continue
except OSError:
return []
return sorted(out)
def resolve_trade_date(data_dir: Path, target: date_cls | None) -> date_cls | None:
"""目标日 → ≤目标日的最近本地交易日。None → None (由 fuyao 取当日)。"""
if target is None:
return None
days = _local_trading_days(data_dir)
if not days:
return target
candidates = [d for d in days if d <= target]
return max(candidates) if candidates else target
def _prev_trading_day(data_dir: Path, d: date_cls) -> date_cls | None:
days = _local_trading_days(data_dir)
earlier = [x for x in days if x < d]
return max(earlier) if earlier else None
def _next_trading_day(data_dir: Path, d: date_cls) -> date_cls | None:
days = _local_trading_days(data_dir)
later = [x for x in days if x > d]
return min(later) if later else None
def _provider():
from app.data_providers import custom as custom_sources
if not custom_sources.is_custom_provider("fuyao"):
return None
return custom_sources.get_provider("fuyao")
def _cache_path(data_dir: Path, d: date_cls) -> Path:
return data_dir / "auction_benchmark" / f"date={d.isoformat()}.json"
def _load_cache(path: Path) -> dict | None:
try:
return json.loads(path.read_text(encoding="utf-8"))
except (OSError, ValueError):
return None
def _store_cache(path: Path, payload: dict) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
tmp = path.with_name(path.name + ".part")
tmp.write_text(json.dumps(payload, ensure_ascii=False), encoding="utf-8")
tmp.replace(path)
def _raw_items(data: dict) -> list[dict]:
items = data.get("item")
return [r for r in items if isinstance(r, dict)] if isinstance(items, list) else []
def _read_kline_closes(data_dir: Path, d: date_cls) -> dict[str, float]:
"""某交易日全市场 {symbol: close}; 无分区/读失败返回空。"""
root = data_dir / "kline_daily" / f"date={d.isoformat()}"
try:
files = sorted(root.glob("*.parquet"))
if not files:
return {}
df = pl.concat([pl.read_parquet(f, columns=["symbol", "close"]) for f in files])
return dict(zip(df["symbol"].to_list(), df["close"].to_list()))
except (OSError, pl.exceptions.PolarsError):
return {}
def _enrich(data_dir: Path, trade_date: date_cls, items: list[dict]) -> list[dict]:
"""名单叠加真实收益: day0_oc (开盘买→收盘卖) / day0_pct (全天) / d1_pct (次日)。
用相邻 kline_daily 分区现算, 不落缓存 — 次日分区晚到时先给 None, 到了自然补上。
"""
root = data_dir / "kline_daily" / f"date={trade_date.isoformat()}"
day0: dict[str, tuple[float, float]] = {} # symbol -> (open, close)
try:
files = sorted(root.glob("*.parquet"))
if files:
df = pl.concat([pl.read_parquet(f, columns=["symbol", "open", "close"]) for f in files])
day0 = dict(zip(df["symbol"].to_list(), zip(df["open"].to_list(), df["close"].to_list())))
except (OSError, pl.exceptions.PolarsError):
day0 = {}
prev = _prev_trading_day(data_dir, trade_date)
prev_close = _read_kline_closes(data_dir, prev) if prev else {}
nxt = _next_trading_day(data_dir, trade_date)
next_close = _read_kline_closes(data_dir, nxt) if nxt else {}
out: list[dict] = []
for r in items:
sym = str(r.get("thscode") or "")
oc = pct = d1 = None
if sym in day0:
o, c = day0[sym]
if o and o > 0 and c is not None:
oc = c / o - 1.0
pc = prev_close.get(sym)
if pc and pc > 0 and c is not None:
pct = c / pc - 1.0
nc = next_close.get(sym)
if c and c > 0 and nc is not None:
d1 = nc / c - 1.0
out.append({
"thscode": sym,
"ticker": r.get("ticker"),
"name": r.get("name"),
"auction_pct": r.get("auction_pct"),
"tags": [str(t) for t in (r.get("tags") or [])],
"day0_oc": oc,
"day0_pct": pct,
"d1_pct": d1,
})
return out
def _base_payload(data_dir: Path, trade_date: date_cls, data: dict) -> dict:
return {
"state": "ok",
"requested_date": trade_date.isoformat(),
"trade_date": str(data.get("date") or trade_date.isoformat()),
"count": len(_raw_items(data)),
"raw_items": _raw_items(data),
}
def get_auction_benchmark(data_dir: Path, target: date_cls | None = None) -> dict:
"""取短线风向标名单 (含当日/次日真实收益)。返回给前端的统一容器。
state: ok (正常) | fallback_prev (目标日拉取失败, 已回退上一期)
| source_unavailable (未配置 fuyao) | no_data (拉取失败)
"""
days = _local_trading_days(data_dir)
if target is None:
target = max(days) if days else None
trade_date = resolve_trade_date(data_dir, target)
today = cn_today()
# 历史日缓存优先 (纯本地, 不触发插件注册表加载)
if trade_date is not None and trade_date < today:
cached = _load_cache(_cache_path(data_dir, trade_date))
if cached is not None:
return _respond(data_dir, trade_date, cached, cached.get("state") or "ok")
provider = _provider()
if provider is None:
return {"state": "source_unavailable"}
from app.plugins.fuyao.client import FuyaoError
explicit = trade_date.isoformat() if trade_date is not None else None
try:
data = provider.short_term_benchmark(explicit)
try:
actual = date_cls.fromisoformat(str(data.get("date")))
except ValueError:
actual = trade_date
base = _base_payload(data_dir, actual, data)
# 历史日不可变 → 落缓存; 当日不缓存 (竞价阶段名单可能变动)
if actual is not None and actual < today:
with contextlib.suppress(OSError):
_store_cache(_cache_path(data_dir, actual), base)
return _respond(data_dir, actual, base, "ok")
except FuyaoError as e:
# 显式日期失败 (非交易日/边界日) → 回退上一交易日一次
if trade_date is not None:
prev = _prev_trading_day(data_dir, trade_date)
if prev is not None:
try:
cached_prev = _load_cache(_cache_path(data_dir, prev))
if cached_prev is not None:
return _respond(data_dir, prev, cached_prev, "fallback_prev",
requested=explicit)
data_prev = provider.short_term_benchmark(prev.isoformat())
base = _base_payload(data_dir, prev, data_prev)
with contextlib.suppress(OSError):
_store_cache(_cache_path(data_dir, prev), base)
return _respond(data_dir, prev, base, "fallback_prev",
requested=explicit)
except FuyaoError:
pass
logger.warning("短线风向标拉取失败: %s", e)
return {"state": "no_data", "message": str(e)}
def _respond(
data_dir: Path,
trade_date: date_cls | None,
base: dict,
state: str,
requested: str | None = None,
) -> dict:
"""缓存/现拉的原始容器 → 叠加收益 enrich 后的前端容器。"""
if trade_date is None:
return {**base, "state": state}
payload = {
"state": state,
"requested_date": requested if requested is not None else base.get("requested_date"),
"trade_date": base.get("trade_date") or trade_date.isoformat(),
"count": base.get("count") or 0,
"items": _enrich(data_dir, trade_date, base.get("raw_items") or []),
}
return payload
def build_recap_context(data_dir: Path) -> str:
"""AI 复盘的盘前风向标摘要段 (纯文本, 失败返回空串不影响复盘)。"""
try:
payload = get_auction_benchmark(data_dir, None)
if payload.get("state") not in ("ok", "fallback_prev"):
return ""
items = payload.get("items") or []
if not items:
return ""
trade_date = payload.get("trade_date") or ""
lines = [f"(数据日期: {trade_date})"]
segs = []
for i in items:
seg = (f"{i.get('name')}({i.get('thscode')}) 竞价{(i.get('auction_pct') or 0):+.2f}%"
f"[{'·'.join(i.get('tags') or [])}]")
if i.get("day0_oc") is not None:
seg += f" → 当日开盘买{i['day0_oc']*100:+.2f}%"
if i.get("d1_pct") is not None:
seg += f", 次日{i['d1_pct']*100:+.2f}%"
segs.append(seg)
lines.append("盘前风向标名单: " + "; ".join(segs))
ocs = [i["day0_oc"] for i in items if i.get("day0_oc") is not None]
if ocs:
lines.append(f"名单当日(开盘买→收盘卖)均值 {sum(ocs)/len(ocs)*100:+.2f}%")
return "\n".join(lines)
except Exception as e: # noqa: BLE001 — 摘要失败不影响复盘主流程
logger.debug("盘前风向标复盘摘要构建失败: %s", e)
return ""
+17 -3
View File
@@ -179,8 +179,9 @@ def _build_emotion_block(overview: dict) -> str:
return "\n".join(lines)
def _build_user_prompt(overview: dict, news: list[dict], focus: str, lhb_context: str = "") -> str:
"""构建用户消息:复盘日期 + 市场数据精简切片 + 龙虎榜(可选) + 新闻 + 关注点。"""
def _build_user_prompt(overview: dict, news: list[dict], focus: str, lhb_context: str = "",
bench_context: str = "") -> str:
"""构建用户消息:复盘日期 + 市场数据精简切片 + 龙虎榜(可选) + 盘前风向标(可选) + 新闻 + 关注点。"""
as_of = overview.get("as_of") or "今日"
parts: list[str] = [
@@ -210,6 +211,14 @@ def _build_user_prompt(overview: dict, news: list[dict], focus: str, lhb_context
lhb_context,
])
# 盘前风向标 (fuyao 竞价筛选名单 + 当日实际表现对照; 失败为空不占段)
if bench_context:
parts.extend([
"",
"## 盘前风向标(竞价)",
bench_context,
])
if news:
news_lines = []
for i, n in enumerate(news[:8], 1):
@@ -312,7 +321,12 @@ async def recap_market_stream(
from app.services import dragon_tiger as dragon_tiger_svc
lhb_ctx = dragon_tiger_svc.build_recap_context(repo.store.data_dir)
user_prompt = _build_user_prompt(overview, news or [], focus, lhb_ctx)
# 盘前风向标摘要 (fuyao 专有): 失败/未配置 → 空串
from app.services import auction_benchmark as auction_benchmark_svc
bench_ctx = auction_benchmark_svc.build_recap_context(repo.store.data_dir)
user_prompt = _build_user_prompt(overview, news or [], focus, lhb_ctx, bench_ctx)
got_content = False
async for delta in stream_ai_text(
[
+178
View File
@@ -0,0 +1,178 @@
"""盘前风向标服务测试 (不依赖真实网络)。
覆盖: 交易日回退、历史日 JSON 缓存命中与落盘、收益 enrich 数学 (当日oc/全天/次日)、
fuyao 未配置降级、目标日失败 fallback_prev、彻底失败 no_data、AI 复盘摘要段。
日期用 2026-08-26/27/28 (写作时为过去交易日), 与仓库既有绝对日期测试风格一致。
"""
from __future__ import annotations
import json
from datetime import date
from pathlib import Path
import polars as pl
import pytest
from app.plugins.fuyao.client import FuyaoError
from app.services import auction_benchmark as ab
def _write_kline(data_dir: Path, day: str, rows: list[tuple[str, float, float]]) -> None:
part = data_dir / "kline_daily" / f"date={day}"
part.mkdir(parents=True, exist_ok=True)
df = pl.DataFrame(
{
"symbol": [r[0] for r in rows],
"open": [r[1] for r in rows],
"close": [r[2] for r in rows],
}
)
df.write_parquet(part / "part-0.parquet")
@pytest.fixture()
def data_dir(tmp_path: Path) -> Path:
for d in ("2026-08-26", "2026-08-27", "2026-08-28"):
(tmp_path / "kline_daily" / f"date={d}").mkdir(parents=True, exist_ok=True)
_write_kline(tmp_path, "2026-08-26", [("600519.SH", 1690.0, 1700.0), ("000858.SZ", 130.0, 131.0)])
_write_kline(tmp_path, "2026-08-27", [("600519.SH", 1717.0, 1734.0), ("000858.SZ", 132.0, 130.0)])
_write_kline(tmp_path, "2026-08-28", [("600519.SH", 1734.0, 1768.68), ("000858.SZ", 129.0, 133.0)])
return tmp_path
class _FakeProvider:
"""记录调用; fail_dates 中的日期抛 FuyaoError。"""
def __init__(self, fail_dates: set[str] | None = None):
self.calls: list[str | None] = []
self.fail_dates = fail_dates or set()
def short_term_benchmark(self, date_iso: str | None) -> dict:
self.calls.append(date_iso)
if date_iso in self.fail_dates:
raise FuyaoError(f"code=3002: {date_iso} 未就绪")
return {
"date": date_iso or "2026-08-28",
"date_ms": 0,
"item": [
{"thscode": "600519.SH", "ticker": "600519", "name": "贵州茅台",
"auction_pct": 1.0, "tags": ["白酒", "超级品牌"]},
{"thscode": "000858.SZ", "ticker": "000858", "name": "五粮液",
"auction_pct": -2.5, "tags": ["白酒"]},
],
}
def _use_provider(monkeypatch, provider) -> _FakeProvider:
monkeypatch.setattr(ab, "_provider", lambda: provider)
return provider
# ---- 交易日解析 ----
def test_resolve_rolls_back_non_trading_day(data_dir):
assert ab.resolve_trade_date(data_dir, date(2026, 8, 30)) == date(2026, 8, 28)
assert ab.resolve_trade_date(data_dir, date(2026, 8, 27)) == date(2026, 8, 27)
# ---- 状态与缓存 ----
def test_source_unavailable_without_fuyao(data_dir, monkeypatch):
monkeypatch.setattr(ab, "_provider", lambda: None)
out = ab.get_auction_benchmark(data_dir, None)
assert out["state"] == "source_unavailable"
def test_fetch_stores_cache_then_hits_cache(data_dir, monkeypatch):
provider = _use_provider(monkeypatch, _FakeProvider())
out = ab.get_auction_benchmark(data_dir, None) # 默认 → 最近分区 08-28
assert out["state"] == "ok" and out["trade_date"] == "2026-08-28"
assert out["count"] == 2 and len(out["items"]) == 2
assert provider.calls == ["2026-08-28"]
assert (data_dir / "auction_benchmark" / "date=2026-08-28.json").exists()
provider.calls.clear()
out2 = ab.get_auction_benchmark(data_dir, date(2026, 8, 30)) # 周日 → 08-28 → 命中缓存
assert out2["state"] == "ok" and out2["trade_date"] == "2026-08-28"
assert provider.calls == []
def test_explicit_history_date_uses_cache(data_dir, monkeypatch):
provider = _use_provider(monkeypatch, _FakeProvider())
ab.get_auction_benchmark(data_dir, date(2026, 8, 27))
assert provider.calls == ["2026-08-27"]
provider.calls.clear()
ab.get_auction_benchmark(data_dir, date(2026, 8, 27))
assert provider.calls == []
def test_failure_falls_back_to_prev(data_dir, monkeypatch):
provider = _use_provider(monkeypatch, _FakeProvider(fail_dates={"2026-08-28"}))
out = ab.get_auction_benchmark(data_dir, date(2026, 8, 28))
assert out["state"] == "fallback_prev"
assert out["trade_date"] == "2026-08-27"
assert out["requested_date"] == "2026-08-28"
# 回退日缓存以 ok 落盘 (不污染直查)
cached = json.loads((data_dir / "auction_benchmark" / "date=2026-08-27.json").read_text(encoding="utf-8"))
assert cached["state"] == "ok"
def test_total_failure_returns_no_data(data_dir, monkeypatch):
_use_provider(monkeypatch, _FakeProvider(fail_dates={"2026-08-28", "2026-08-27"}))
out = ab.get_auction_benchmark(data_dir, date(2026, 8, 28))
assert out["state"] == "no_data"
assert "2026-08-28" in out.get("message", "")
def test_corrupt_cache_refetches(data_dir, monkeypatch):
provider = _use_provider(monkeypatch, _FakeProvider())
cache = data_dir / "auction_benchmark" / "date=2026-08-28.json"
cache.parent.mkdir(parents=True, exist_ok=True)
cache.write_text("{broken json", encoding="utf-8")
out = ab.get_auction_benchmark(data_dir, date(2026, 8, 28))
assert out["state"] == "ok"
assert provider.calls # 缓存损坏 → 重新拉取
# ---- 收益 enrich ----
def test_enrich_math_with_local_kline(data_dir, monkeypatch):
# 显式查 08-27: prev=08-26, next=08-28
_use_provider(monkeypatch, _FakeProvider())
out = ab.get_auction_benchmark(data_dir, date(2026, 8, 27))
by = {i["thscode"]: i for i in out["items"]}
mt = by["600519.SH"]
# day0_oc = 1734/1717-1; day0_pct = 1734/1700-1; d1 = 1768.68/1734-1
assert mt["day0_oc"] == pytest.approx(1734.0 / 1717.0 - 1)
assert mt["day0_pct"] == pytest.approx(1734.0 / 1700.0 - 1)
assert mt["d1_pct"] == pytest.approx(1768.68 / 1734.0 - 1)
wly = by["000858.SZ"]
assert wly["day0_oc"] == pytest.approx(130.0 / 132.0 - 1)
assert wly["d1_pct"] == pytest.approx(133.0 / 130.0 - 1)
# 原始字段透传
assert mt["auction_pct"] == 1.0 and mt["tags"] == ["白酒", "超级品牌"]
def test_enrich_missing_kline_gives_none(data_dir, monkeypatch):
# 最新分区 08-28 无次日 → d1_pct=None; kline 行存在则 oc/pct 正常
_use_provider(monkeypatch, _FakeProvider())
out = ab.get_auction_benchmark(data_dir, None)
for i in out["items"]:
assert i["d1_pct"] is None
assert i["day0_oc"] is not None
# ---- AI 复盘摘要 ----
def test_build_recap_context_contains_summary(data_dir, monkeypatch):
_use_provider(monkeypatch, _FakeProvider())
ctx = ab.build_recap_context(data_dir)
assert "盘前风向标名单" in ctx and "贵州茅台" in ctx
assert "白酒" in ctx # 概念标签
assert "当日" in ctx # 收益对照
def test_build_recap_context_empty_without_source(data_dir, monkeypatch):
monkeypatch.setattr(ab, "_provider", lambda: None)
assert ab.build_recap_context(data_dir) == ""