diff --git a/backend/app/plugins/fuyao/client.py b/backend/app/plugins/fuyao/client.py index d206b5e..f13693e 100644 --- a/backend/app/plugins/fuyao/client.py +++ b/backend/app/plugins/fuyao/client.py @@ -1,12 +1,17 @@ """扶摇(同花顺金融数据 API) HTTP 客户端。 -职责: 认证、统一信封解包、分页拉取快照。不知道 provider / services 层。 +职责: 认证、统一信封解包、快照分页、单标的日K、市场 dump 下载。不知道 provider / services 层。 文档: https://fuyao.aicubes.cn/docs — REST + X-api-key, 响应信封 {code, message, data}。 + +时间字段口径: 所有 *ms 字段(含 start/end 入参与 date_ms/ex_date_ms 出参)均为 +北京时间零点对应的 epoch ms(= UTC 前一日 16:00), 由 provider 层统一 +8h 换算。 """ + from __future__ import annotations import logging import time +from pathlib import Path import httpx @@ -60,7 +65,9 @@ class FuyaoClient: return payload.get("data") or {} # ---- 快照 ---- - def snapshot_page(self, limit: int = _SNAPSHOT_PAGE_SIZE, offset: int = 0) -> tuple[list[dict], int]: + def snapshot_page( + self, limit: int = _SNAPSHOT_PAGE_SIZE, offset: int = 0 + ) -> tuple[list[dict], int]: """拉取一页 A 股全市场快照。返回 (rows, total), total 为全市场总数。 实测响应(2026-08): data={timestamp, total, item}; 官方文档示例为 @@ -107,3 +114,131 @@ class FuyaoClient: if not out: raise FuyaoError("全市场快照为空") return out, server_ts + + # ---- 历史日K ---- + def historical_kline( + self, thscode: str, start_ms: int, end_ms: int, adjust: str = "none" + ) -> list[dict]: + """单标的日K(interval=1d 固定)。单次窗口 ≤10 年, 超出由调用方分片。 + + adjust 必须显式传 "none" 取原始价 — 服务端默认是 forward(前复权), + 官方前复权序列事件间存在逐日漂移, 项目内禁止使用。 + 返回 data.item 原始行: {date_ms, open_price, high_price, low_price, + close_price, volume(股), turnover(元)}。 + """ + data = self._get( + "/api/a-share/prices/historical", + { + "thscode": thscode, + "interval": "1d", + "adjust": adjust, + "start": int(start_ms), + "end": int(end_ms), + }, + ) + rows = data.get("item") + return rows if isinstance(rows, list) else [] + + # ---- 财务 ---- + # 端点均单标的(thscode 不接受逗号)。取数模式二选一: limit=最近N期 或 start/end 区间, + # 这里只用 limit。period=quarterly 覆盖每个季度末(含年报期), 与项目"各报告期累积"口径一致。 + _STATEMENT_ENDPOINTS = { + "income": "income-statements", + "balance_sheet": "balance-sheets", + "cash_flow": "cash-flow-statements", + } + + def financial_statements( + self, stmt: str, thscode: str, limit: int = 1 + ) -> list[dict]: + """单标的财务报表多期序列。stmt: income | balance_sheet | cash_flow。 + + 返回 data.item 原始行: 共有元数据(thscode/period/fiscal_year/fiscal_period/ + report_date_ms/period_end_ms/currency) + 各表字段。行内 null 表示该期未披露。 + """ + endpoint = self._STATEMENT_ENDPOINTS.get(stmt) + if endpoint is None: + raise FuyaoError(f"未知财务报表类型: {stmt}") + data = self._get( + f"/api/a-share/financials/{endpoint}", + {"thscode": thscode, "period": "quarterly", "limit": max(1, min(20, limit))}, + ) + rows = data.get("item") + return rows if isinstance(rows, list) else [] + + def financial_indicators(self, thscode: str, report: str) -> list[dict]: + """单标的单报告期财务指标(report 格式 yyyy-N, N=1..4 对应一季报..年报)。 + + 返回 data.abilities 原始列表 [{ability, indicators: [{index_id, value}]}]; + value 为保留原始精度的数值字符串(百分制指标即百分点数), 缺失为 null。 + 未披露报告期实测返回 code=5003(文档写 3002, 以实测为准) → 经 _get 抛 FuyaoError, + 由调用方按"该期无数据"处理。 + """ + data = self._get( + "/api/a-share/financials/indicators", + {"thscode": thscode, "report": report}, + ) + abilities = data.get("abilities") + return abilities if isinstance(abilities, list) else [] + + def valuations_snapshot(self, thscodes: list[str]) -> list[dict]: + """批量估值快照(pe_ttm/pe_mrq/pb_mrq/ps_ttm/pcf_ttm), 数值为最新口径。 + + 服务端单次上限 100 只(超出 code=1003), 分批由调用方负责。 + 返回 data.item 原始行。 + """ + data = self._get( + "/api/a-share/valuations/snapshot", + {"thscodes": ",".join(thscodes[:100])}, + ) + rows = data.get("item") + return rows if isinstance(rows, list) else [] + + def price_snapshot_batch(self, thscodes: list[str]) -> list[dict]: + """按 thscodes 批量行情快照(最新价等), 用于估值推导的分母。 + + 与全市场分页快照同一端点; thscodes 显式传入时不分页。 + 返回 data.item 原始行。 + """ + data = self._get( + "/api/a-share/prices/snapshot", + {"thscodes": ",".join(thscodes[:100])}, + ) + rows = data.get("item") + return rows if isinstance(rows, list) else [] + + # ---- 市场 dump ---- + def dump_download_url(self, dump_kind: str) -> dict: + """获取 dump 预签名下载信息(约 300s 有效)。 + + dump_kind: adjustment-factors | daily-k-10d | daily-k。 + 返回 {presigned_url, presigned_url_expires_at, expires_in_seconds}; + release 版本号(如 20260828)嵌在 presigned_url 的 releases// 路径中, + 供调用方做缓存版本管理。 + """ + return self._get(f"/api/dump/market-dumps/{dump_kind}/download-url", {}) + + def download_dump(self, dump_kind: str, dest: Path) -> Path: + """下载 dump 到 dest(先写 .part 临时文件, 成功后原子改名)。失败抛 FuyaoError。 + + 预签名 URL 指向对象存储, 请求不得携带 X-api-key 头 → 用独立裸请求, + 不经过持有认证头的 self._http。 + """ + url = str(self.dump_download_url(dump_kind).get("presigned_url") or "") + if not url: + raise FuyaoError(f"dump {dump_kind} 未返回预签名 URL") + dest.parent.mkdir(parents=True, exist_ok=True) + tmp = dest.with_name(dest.name + ".part") + try: + with httpx.stream("GET", url, timeout=120.0, follow_redirects=True) as resp: + if resp.status_code != 200: + raise FuyaoError(f"dump {dump_kind} 下载失败 HTTP {resp.status_code}") + with open(tmp, "wb") as fh: + for chunk in resp.iter_bytes(1 << 20): + fh.write(chunk) + tmp.replace(dest) + except httpx.HTTPError as e: + raise FuyaoError(f"dump {dump_kind} 下载网络失败: {e}") from e + finally: + tmp.unlink(missing_ok=True) + return dest diff --git a/backend/app/plugins/fuyao/plugin.yaml b/backend/app/plugins/fuyao/plugin.yaml index 7e11613..cc063f7 100644 --- a/backend/app/plugins/fuyao/plugin.yaml +++ b/backend/app/plugins/fuyao/plugin.yaml @@ -7,8 +7,8 @@ display_name: "fuyao" runtime: none entry: app.plugins.fuyao.provider:FuyaoProvider check: app.plugins.fuyao.provider:availability -datasets: [realtime] +datasets: [realtime, daily, adj_factor, financial] api_key_env: FUYAO_API_KEY # 声明后设置页提供 Key 输入框(先探后存, secrets.json 优先) -description: "同花顺官方 REST 数据 API。当前提供 A 股全市场实时快照(分页拉取);日K/分钟/财务未接入,自动回退 TickFlow。" +description: "同花顺官方 REST 数据 API。A 股实时快照、日K(原始价: 近端走 10d dump, 深窗口走 10 年全量 dump 一次下载秒级筛选, 兜底单标的接口)、除权因子(事件 dump + 本地日K dump 配价推导, 涨跌停自检, 秒级)、财务(利润表/资产负债表/现金流量表/指标, 字段映射至项目口径, bps 由估值反推; 股本无上游接口, 指标历史期建议切回 TickFlow 补齐);分钟未接入,自动回退 TickFlow。" install_hint: "点击卡片中的输入框配置 API Key(https://fuyao.aicubes.cn 申请),或在 .env 中配置 FUYAO_API_KEY" homepage: "https://fuyao.aicubes.cn" diff --git a/backend/app/plugins/fuyao/provider.py b/backend/app/plugins/fuyao/provider.py index f0c7379..93def17 100644 --- a/backend/app/plugins/fuyao/provider.py +++ b/backend/app/plugins/fuyao/provider.py @@ -3,35 +3,77 @@ 方法签名对齐 custom.GenericHTTPProvider(service 分流点按这套签名调用), 注入 custom loader 注册表后, 各 service 无需改动即可路由到本 provider。 -当前实现数据集: realtime (A 股全市场快照, 分页)。 -未声明 daily / minute / financial → provider_has_dataset 为 False, 自动回退 tickflow。 +实现数据集: + - realtime A 股全市场快照 (分页) + - daily A 股日K, 原始价; 近端窗口走 daily-k-10d 全市场 dump(1 次请求), + 深窗口走单标的 historical 接口(≤10 年/次自动分片) + - adj_factor A 股除权因子; adjustment-factors 事件 dump + 自家原始日K前收盘, + 按交易所公式推导单事件比值, 涨跌停自检 + - financial 财务五表(股本除外): 三表多期序列 + 指标单期, 字段映射为 TickFlow + canonical 列名, 扶摇独有字段原名透传为扩展列; bps 由估值 pb_mrq + 反推; shares 无上游接口恒空 +未声明 minute → provider_has_dataset 为 False, 自动回退 tickflow。 -单位口径 (CONTRIBUTING §3.1, 不可凭字段名推断): +单位与口径 (CONTRIBUTING §3.1, 不可凭字段名推断): - 扶摇 price_change_ratio_pct 为百分数数值 (1.74 = +1.74%), 本项目 realtime change_pct 契约为小数制 (0.0174 = 1.74%) → 此处显式 / 100。 - - volume 单位股、turnover 单位元, 与内部契约一致, 直接透传。 + - 扶摇 volume 单位为股, 本项目日K/实时契约均为手 → 统一 floor(股/100)。 + - turnover 单位元, 与内部一致, 直接透传。 + - 日K取数 adjust=none 锁定: 官方 forward 序列事件间有逐日漂移(2026-08 实测), + 项目内前复权一律由 indicators.pipeline 用本地因子计算。 + - ex_factor 为单事件比值(非累积), 累积链由 pipeline._apply_adj_factor 构建。 """ + from __future__ import annotations +import calendar import contextlib import logging +import math +import re import time +from collections.abc import Callable from dataclasses import dataclass, field +from datetime import UTC, date, datetime, timedelta +from pathlib import Path +import polars as pl + +from app.data_providers.normalizer import DAILY_COLS, normalize_daily +from app.indicators.pipeline import filter_halt_days from app.plugins.fuyao import client as fuyao_client from app.plugins.fuyao.client import FuyaoClient, FuyaoError logger = logging.getLogger(__name__) # 只声明真实提供的数据集; 其余数据集 provider_has_dataset 返回 False → 回退 tickflow -_DATASETS = ("realtime",) +_DATASETS = ("realtime", "daily", "adj_factor", "financial") API_KEY_ENV = "FUYAO_API_KEY" SECRETS_FIELD = "fuyao_api_key" # UI 配置的 Key 存 secrets.json, 优先级高于 .env +# 扶摇 *ms 时间字段为北京时间零点(= UTC 前一日 16:00), +8h 后按 UTC 解析即得交易日 +_SH_MS = 28_800_000 +_HIST_MAX_SPAN_MS = 3650 * 86_400_000 # historical 单次窗口上限 10 年, 超出由本层分片 +_HIST_INTERVAL_S = 0.12 # 单标的请求节流(实测 200+ 连发未触发 4001 限频) +_FINANCIAL_HISTORY_PERIODS = 8 # 财务首装全量历史: 最近 8 期季报(约 2 年) +_VALUATION_BATCH = 100 # 估值/价格快照端点单次上限 100 只 +# 项目财务表名 → 扶摇报表端点名 +_STATEMENT_ENDPOINTS = { + "income": "income-statements", + "balance_sheet": "balance-sheets", + "cash_flow": "cash-flow-statements", +} +_ADJ_DUMP_KIND = "adjustment-factors" +_DAILY10_DUMP_KIND = "daily-k-10d" +_DAILY_DUMP_KIND = "daily-k" # 10 年全量日K dump(约 172MB), 深窗口一次下载覆盖全市场 +_RECENT_DUMP_DAYS = 12 # 窗口跨度 ≤ 此天数时优先走 10d dump(覆盖 ≈10 个交易日) +_PREV_CLOSE_BACKDAYS = 30 # 推导因子时向前找"除权日前收盘"的回看天数(容忍长期停牌) + def get_api_key() -> str: from app import secrets_store + return secrets_store.get_env_backed_secret(SECRETS_FIELD, API_KEY_ENV) @@ -39,7 +81,8 @@ def availability() -> tuple[bool, str]: """loader 启动自检: API Key 已配置(secrets.json 或 .env)才注册为可切换数据源。不抛异常。""" if get_api_key(): return True, "ok" - return False, f"未配置 {API_KEY_ENV}(可在设置页数据源卡片中直接填写)" + # 状态行会拼在「未配置」标签之后, 文案不再重复"未配置"字样 + return False, f"缺少 API Key(可在下方输入框直接填写,或配置环境变量 {API_KEY_ENV})" def probe_api_key(api_key: str) -> tuple[bool, str]: @@ -85,6 +128,128 @@ def _first(row: dict, *names: str): return None +def _date_of_ms(value) -> date | None: + """扶摇 *ms(上海零点) → 交易日。None/非法值返回 None, 不伪造。""" + if value is None: + return None + try: + ms = int(value) + except (TypeError, ValueError): + return None + return datetime.fromtimestamp((ms + _SH_MS) // 1000, tz=UTC).date() + + +def _ms_of_date(d: date) -> int: + """交易日 → 扶摇 start/end 入参口径的 ms(该日上海零点的 epoch ms, 不依赖本机时区)。""" + return (calendar.timegm(d.timetuple()) - 28_800) * 1000 + + +def _iso_of_ms(value) -> str | None: + """扶摇 *ms → ISO 日期字符串(项目财务表 period_end/announce_date 的存储口径)。""" + d = _date_of_ms(value) + return d.isoformat() if d is not None else None + + +def _report_quarter(fiscal_period) -> int | None: + """fiscal_period(Q1..Q4/FY) → 指标接口 report 参数的季号 N(1..4)。""" + if not fiscal_period: + return None + text = str(fiscal_period).strip().upper() + if text == "FY": + return 4 + try: + return int(text.lstrip("Q")) + except ValueError: + return None + + +def _ref_price( + prev_close: float, dividend: float, bonus: float, allot: float, allot_price: float +) -> float | None: + """交易所除权参考价: (P - D + AR·AP) / (1 + S + AR), 四舍五入(half-up)保留 2 位。 + + half-up 是交易所口径; 银行家舍入会让约半数事件在第 2 位小数上偏离(对拍实证)。 + """ + denom = 1.0 + bonus + allot + if denom <= 0: + return None + x = (prev_close - dividend + allot * allot_price) / denom + return math.floor(x * 100 + 0.5) / 100 + + +def _price_limit(symbol: str) -> float: + """按代码前缀给涨跌停幅度(自检容差用): 创业板/科创板 20%, 北交所 30%, 主板 10%。""" + code = symbol.split(".")[0] + if code.startswith(("300", "301", "688", "689")): + return 0.20 + if code.startswith(("8", "4", "92")): + return 0.30 + return 0.10 + + +def _release_of(url: str) -> str: + """从预签名 URL 提取 release 版本号(releases// 路径), 提不到返回 unknown。""" + m = re.search(r"releases/(\d+)/", url or "") + return m.group(1) if m else "unknown" + + +def _cache_dir() -> Path: + from app.config import settings + + d = settings.data_dir / "cache" / "fuyao" + d.mkdir(parents=True, exist_ok=True) + return d + + +def _kline_rows(symbol: str, bars: list[dict]) -> list[dict]: + """historical/dump 原始行(价格元, volume 股) → 内部行(volume 手)。""" + out = [] + for b in bars: + v = _to_float(b.get("volume")) + out.append( + { + "symbol": symbol, + "date": _date_of_ms(b.get("date_ms")), + "open": _to_float(b.get("open_price")), + "high": _to_float(b.get("high_price")), + "low": _to_float(b.get("low_price")), + "close": _to_float(b.get("close_price")), + "volume": math.floor(v / 100.0) if v is not None else None, + "amount": _to_float(b.get("turnover")), + } + ) + return out + + +def _tail_ok(end_d: date, covered_max: date) -> bool: + """请求终点是否被覆盖到 covered_max: 周末/节假日的自然缺口(≤3 天)不算缺失。""" + if end_d <= covered_max: + return True + return (end_d - covered_max).days <= 3 and end_d.weekday() >= 5 + + +def _dump_covers(dump: pl.DataFrame, start_d: date, end_d: date) -> bool: + """dump 日期范围是否覆盖请求窗口: 起点必须落在 dump 内; 终点允许周末自然缺口。""" + if dump.is_empty() or "date_ms" not in dump.columns: + return False + dates = pl.from_epoch(dump["date_ms"].cast(pl.Int64) + _SH_MS, time_unit="ms").dt.date() + dmin, dmax = dates.min(), dates.max() + return start_d >= dmin and _tail_ok(end_d, dmax) + + +def _dump_date_range(path: Path) -> tuple[date | None, date | None]: + """lazy 读 parquet 的 date_ms 边界(走元数据/少量行组, 不整读大文件)。""" + row = ( + pl.scan_parquet(path) + .select( + pl.from_epoch(pl.col("date_ms").min() + _SH_MS, time_unit="ms").dt.date().alias("dmin"), + pl.from_epoch(pl.col("date_ms").max() + _SH_MS, time_unit="ms").dt.date().alias("dmax"), + ) + .collect() + ) + return row["dmin"][0], row["dmax"][0] + + def _map_snapshot_row(row: dict, fetched_ms: int) -> dict | None: """扶摇快照行 → 内部 realtime record。字段缺失时按依赖推导, 不伪造数据。 @@ -108,6 +273,8 @@ def _map_snapshot_row(row: dict, fetched_ms: int) -> dict | None: # 与 quote_service 的推导同口径: 小数制, 不乘 100 change_pct = change_amount / prev + volume = _to_float(row.get("volume")) + return { "symbol": symbol, "name": row.get("name"), # 快照无名称, 由下游维表关联 @@ -116,11 +283,11 @@ def _map_snapshot_row(row: dict, fetched_ms: int) -> dict | None: "open": _to_float(row.get("open_price")), "high": _to_float(_first(row, "high_price", "highest_price")), "low": _to_float(_first(row, "low_price", "lowest_price")), - "volume": _to_float(row.get("volume")), + "volume": math.floor(volume / 100.0) if volume is not None else None, # 股 → 手 "amount": _to_float(row.get("turnover")), "change_pct": change_pct, "change_amount": change_amount, - "amplitude": None, # 快照未提供, 不启发式计算 + "amplitude": None, # 快照未提供, 不启发式计算 "turnover_rate": None, # 需股本口径 (§3.4), 交给 enriched 管道用历史股本计算 "timestamp": fetched_ms, "session": None, @@ -136,18 +303,76 @@ class FuyaoProvider: def __init__(self) -> None: self.config = _FuyaoConfig() self._client: FuyaoClient | None = None + self._dump_memo: dict[str, pl.DataFrame] = {} + self._dump_path_memo: dict[str, Path] = {} def close(self) -> None: # loader.load_all 重建注册表时会对每个 provider 调 close if self._client is not None: with contextlib.suppress(Exception): self._client.close() self._client = None + self._dump_memo.clear() + self._dump_path_memo.clear() def _get_client(self) -> FuyaoClient: if self._client is None: self._client = fuyao_client.FuyaoClient(api_key=get_api_key()) return self._client + # ---- dump 缓存 ---- + def _ensure_dump_path(self, dump_kind: str, cache_prefix: str) -> Path: + """确保最新 release 的 dump 已落盘, 返回缓存路径(大文件不整读进内存)。 + + release 号取自预签名 URL 的 releases// 路径; 新 release 落盘后清理旧版缓存。 + """ + memo = self._dump_path_memo.get(dump_kind) + if memo is not None and memo.exists(): + return memo + client = self._get_client() + info = client.dump_download_url(dump_kind) + release = _release_of(str(info.get("presigned_url") or "")) + dest = _cache_dir() / f"{cache_prefix}__{release}.parquet" + if not dest.exists(): + client.download_dump(dump_kind, dest) + for old in dest.parent.glob(f"{cache_prefix}__*.parquet"): + if old.name != dest.name: + old.unlink(missing_ok=True) + logger.info("扶摇 dump %s(release %s)已下载: %s", dump_kind, release, dest.name) + self._dump_path_memo[dump_kind] = dest + return dest + + def _ensure_dump(self, dump_kind: str, cache_prefix: str) -> pl.DataFrame: + """小体量 dump(快照 10d / 因子)整读 + 进程内 memo, 避免重复打接口/读盘。""" + memo = self._dump_memo.get(dump_kind) + if memo is not None: + return memo + df = pl.read_parquet(self._ensure_dump_path(dump_kind, cache_prefix)) + self._dump_memo[dump_kind] = df + return df + + def _ensure_daily_big_dump(self, start_d: date) -> Path | None: + """10 年全量日K dump(约 172MB)。只要求覆盖窗口起点; 末端缺口由 10d dump 补。 + + 已有缓存覆盖起点就直接复用, 不追新 release(避免深窗口高频触发时日日重下 + 172MB — 旧 release 的中段历史不会变, 尾部新鲜度交给 1MB 的 10d dump)。 + """ + for f in sorted(_cache_dir().glob("daily_k__*.parquet"), reverse=True): + try: + dmin, _ = _dump_date_range(f) + except Exception: # 缓存损坏不致命, 换下一个/重拉 + continue + if dmin is not None and start_d >= dmin: + return f + try: + path = self._ensure_dump_path(_DAILY_DUMP_KIND, "daily_k") + except FuyaoError as e: + logger.warning("扶摇 10 年 dump 不可用, 回退单标的接口: %s", e) + return None + dmin, _ = _dump_date_range(path) + if dmin is None or start_d < dmin: + return None # 窗口比 10 年更早 → 单标的兜底 + return path + # ---- realtime ---- def get_realtime(self) -> list[dict]: """全市场实时快照 → 内部 realtime records。失败软返回空列表(不阻断轮询)。""" @@ -175,11 +400,636 @@ class FuyaoProvider: logger.info("扶摇实时行情拉取完成: %d 条(丢弃 %d 行)", len(records), dropped) return records + # ---- daily ---- + def get_daily( + self, + symbols: list[str], + start_time: datetime | None, + end_time: datetime | None, + asset_type: str = "stock", + on_chunk_done: Callable[[int, int], None] | None = None, + ) -> pl.DataFrame: + """A 股日K → 内部契约。原始价(adjust=none 锁定)、volume 股→手(floor /100)。 + + 取数三档(逐级降级, 全部失败才空手而归): + - 近端窗口(跨度 ≤12 天): daily-k-10d dump(约 1MB), 1 次请求覆盖全部标的; + - 深窗口: daily-k 10 年全量 dump(约 172MB, 缓存覆盖起点即复用, 不追新 release), + 末端缺口由 10d dump 补尾 — 全市场深回填从"逐票 34 分钟"降为"一次下载+秒级筛选"; + - 兜底: 单标的 historical 接口(窗口早于 dump 覆盖 / dump 不可用; 10 年自动分片, + 逐标的节流 + 进度回调)。 + """ + if not symbols or asset_type != "stock": + return pl.DataFrame() + end_dt = end_time or datetime.now() + start_dt = start_time or (end_dt - timedelta(days=365)) + start_d, end_d = start_dt.date(), end_dt.date() + + if (end_d - start_d).days <= _RECENT_DUMP_DAYS: + try: + dump = self._ensure_dump(_DAILY10_DUMP_KIND, "daily_k_10d") + if _dump_covers(dump, start_d, end_d): + df = self._daily_from_dump(dump, set(symbols), start_d, end_d) + if on_chunk_done: + on_chunk_done(1, 1) + logger.info("扶摇日K(10d dump)完成: %d 行 [%s ~ %s]", df.height, start_d, end_d) + return df + logger.info("扶摇 10d dump 未覆盖窗口 [%s ~ %s], 尝试 10 年 dump", start_d, end_d) + except FuyaoError as e: + logger.warning("扶摇日K 10d dump 不可用, 尝试 10 年 dump: %s", e) + + df = self._daily_from_big_dump(set(symbols), start_d, end_d) + if df is not None: + if on_chunk_done: + on_chunk_done(1, 1) + logger.info("扶摇日K(10 年 dump)完成: %d 行 [%s ~ %s]", df.height, start_d, end_d) + return df + df = self._daily_from_api(symbols, start_d, end_d, on_chunk_done) + logger.info("扶摇日K(单标的接口)完成: %d 行 [%s ~ %s]", df.height, start_d, end_d) + return df + + def _daily_from_dump( + self, dump: pl.DataFrame, symset: set[str], start_d: date, end_d: date + ) -> pl.DataFrame: + df = dump.with_columns( + pl.from_epoch(pl.col("date_ms") + _SH_MS, time_unit="ms").dt.date().alias("date") + ) + return self._map_daily_dump(df, symset, start_d, end_d) + + def _map_daily_dump( + self, df: pl.DataFrame, symset: set[str], start_d: date, end_d: date + ) -> pl.DataFrame: + """dump 原始行 → 内部日K契约(股→手、字段重命名、停牌过滤)。""" + if "adjusted" in df.columns: + df = df.filter(pl.col("adjusted") == "none") # 防御: dump 变为复权口径时拒绝落库 + out = ( + df.filter( + (pl.col("date") >= start_d) + & (pl.col("date") <= end_d) + & pl.col("thscode").is_in(sorted(symset)) + ) + .with_columns( + (pl.col("volume") / 100.0).floor().alias("volume"), # 股 → 手 + pl.col("thscode").alias("symbol"), + ) + .rename( + { + "open_price": "open", + "high_price": "high", + "low_price": "low", + "close_price": "close", + "turnover": "amount", + } + ) + ) + out = filter_halt_days(out) + cols = [c for c in DAILY_COLS if c in out.columns] + return out.select(cols).sort(["symbol", "date"]) if not out.is_empty() else out.select(cols) + + def _daily_from_big_dump( + self, symset: set[str], start_d: date, end_d: date + ) -> pl.DataFrame | None: + """深窗口主路径: 10 年全量 dump(lazy 按需筛) + 必要时 10d dump 补尾。 + + 覆盖不了(窗口早于 10 年 / dump 拉取失败)返回 None, 由调用方走单标的接口。 + """ + path = self._ensure_daily_big_dump(start_d) + if path is None: + return None + _, dmax = _dump_date_range(path) + big_hi = min(end_d, dmax) + # 窗口/标的过滤下推到 lazy 计划, 只物化需要的行(全量 10 年 ≈ 13.6M 行) + window = ( + pl.scan_parquet(path) + .with_columns( + pl.from_epoch(pl.col("date_ms") + _SH_MS, time_unit="ms").dt.date().alias("date") + ) + .filter( + (pl.col("date") >= start_d) + & (pl.col("date") <= big_hi) + & pl.col("thscode").is_in(sorted(symset)) + ) + .collect() + ) + parts = [self._map_daily_dump(window, symset, start_d, big_hi)] + if not _tail_ok(end_d, dmax): + # 末端缺口(如 10 年 dump 是旧 release, end 是最近交易日): 10d dump 补尾 + try: + ten = self._ensure_dump(_DAILY10_DUMP_KIND, "daily_k_10d") + ten_dates = pl.from_epoch( + ten["date_ms"].cast(pl.Int64) + _SH_MS, time_unit="ms" + ).dt.date() + ten_min, ten_max = ten_dates.min(), ten_dates.max() + if ( + ten_min is not None + and ten_min <= dmax + timedelta(days=1) + and _tail_ok(end_d, ten_max) + ): + tail_start = max(start_d, dmax + timedelta(days=1)) + parts.append(self._daily_from_dump(ten, symset, tail_start, end_d)) + else: + return None # 中段或尾部仍有缺口 → 单标的兜底, 不交缺口数据 + except FuyaoError as e: + logger.warning("扶摇 10d dump 补尾失败: %s", e) + return None + non_empty = [p for p in parts if not p.is_empty()] + if not non_empty: + return pl.DataFrame() + out = pl.concat(non_empty, how="vertical_relaxed") + return out.unique(subset=["symbol", "date"], keep="last").sort(["symbol", "date"]) + + def _daily_from_api( + self, + symbols: list[str], + start_d: date, + end_d: date, + on_chunk_done: Callable[[int, int], None] | None, + ) -> pl.DataFrame: + frames: list[pl.DataFrame] = [] + for i, sym in enumerate(symbols): + rows = self._historical_bars(sym, start_d, end_d) + time.sleep(_HIST_INTERVAL_S) + if rows: + df = normalize_daily(_kline_rows(sym, rows), default_symbol=sym, source=self.name) + if not df.is_empty(): + frames.append(df) + if on_chunk_done: + on_chunk_done(i + 1, len(symbols)) + return pl.concat(frames, how="diagonal_relaxed") if frames else pl.DataFrame() + + def _historical_bars(self, symbol: str, start_d: date, end_d: date) -> list[dict]: + """按 ≤10 年窗口分片拉取单标的原始日K。中途失败软返回已得行, 不抛出。""" + out: list[dict] = [] + s = _ms_of_date(start_d) + e = _ms_of_date(end_d) + while s <= e: + chunk_end = min(e, s + _HIST_MAX_SPAN_MS - 1) + try: + out.extend(self._get_client().historical_kline(symbol, s, chunk_end, adjust="none")) + except FuyaoError as err: + logger.warning("扶摇日K拉取失败 %s [%s ~ %s]: %s", symbol, start_d, end_d, err) + break + if s + _HIST_MAX_SPAN_MS <= e: + time.sleep(_HIST_INTERVAL_S) + s = chunk_end + 1 + return out + + # ---- adj_factor ---- + def get_adj_factors( + self, + symbols: list[str], + start_time: datetime | None, + end_time: datetime | None, + asset_type: str = "stock", + on_chunk_done: Callable[[int, int], None] | None = None, + ) -> pl.DataFrame: + """A 股除权因子 → 内部契约(symbol/trade_date/ex_factor, 单事件比值非累积)。 + + 数据链: adjustment-factors 全量事件 dump → 清洗(时区 +8h / 滤全零 / 同日成分合并) + → 前收盘价从本地日K dump 一次取齐(缺价标的回退单标的接口) → 交易所公式推导 + → 涨跌停自检剔除异常。 + 未来已公告事件无前收盘, 留给滚动增量窗口(15 天)自然补上。 + """ + schema = {"symbol": pl.String, "trade_date": pl.Date, "ex_factor": pl.Float64} + if not symbols or asset_type != "stock": + return pl.DataFrame(schema=schema) + try: + events = self._load_adj_events(set(symbols), start_time, end_time) + except FuyaoError as e: + logger.warning("扶摇除权因子 dump 加载失败: %s", e) + return pl.DataFrame(schema=schema) + if events.is_empty(): + return pl.DataFrame(schema=schema) + + out_rows: list[dict] = [] + syms = sorted(events["symbol"].unique().to_list()) + # 前收盘价优先从本地日K dump 一次取齐(全市场配价从逐标的 ~13 分钟降为秒级), + # 大 dump 不可用或个别标的缺价时回退单标的接口(节流保留)。 + bounds = events.group_by("symbol").agg( + pl.col("ex_date").min().alias("first_ex"), + pl.col("ex_date").max().alias("last_ex"), + ) + closes_by_sym = self._closes_from_dumps(bounds) + fallbacks = 0 + for i, sym in enumerate(syms): + evs = events.filter(pl.col("symbol") == sym).sort("ex_date") + first_ex: date = evs["ex_date"][0] + last_ex: date = evs["ex_date"][-1] + closes = (closes_by_sym or {}).get(sym) + if not closes or min(closes) > first_ex: + # dump 无该标的 / 覆盖不到首个事件前 → 单标的接口兜底 + closes = self._fetch_closes( + sym, first_ex - timedelta(days=_PREV_CLOSE_BACKDAYS), last_ex + ) + time.sleep(_HIST_INTERVAL_S) + fallbacks += 1 + if not closes: + logger.warning("扶摇除权因子: %s 原始日K为空, 跳过其 %d 个事件", sym, evs.height) + continue + days = sorted(closes) + for ev in evs.iter_rows(named=True): + exd: date = ev["ex_date"] + prev_days = [d for d in days if d < exd] + if not prev_days: + continue # 日K窗口未覆盖的远古事件(如 10 年分片边界之外) + p = closes[prev_days[-1]] + ref = _ref_price(p, ev["dividend"], ev["bonus"], ev["allot"], ev["allot_price"]) + if ref is None or ref <= 0: + logger.warning( + "扶摇除权因子: %s %s 参考价无法计算(P=%s D=%s S=%s AR=%s AP=%s), 跳过", + sym, + exd, + p, + ev["dividend"], + ev["bonus"], + ev["allot"], + ev["allot_price"], + ) + continue + factor = p / ref + ex_days = [d for d in days if d >= exd] + if ex_days: # 涨跌停自检: 错误因子会令复权后除权日涨跌幅超出板限 + ret = closes[ex_days[0]] / (p / factor) - 1.0 + if abs(ret) > _price_limit(sym) + 0.02: + logger.warning( + "扶摇除权因子自检剔除 %s %s: 复权后除权日涨跌幅 %.1f%% 超出涨跌停", + sym, + exd, + ret * 100, + ) + continue + out_rows.append({"symbol": sym, "trade_date": exd, "ex_factor": factor}) + if on_chunk_done: + on_chunk_done(i + 1, len(syms)) + if closes_by_sym is not None: + logger.info( + "扶摇除权因子: 本地 dump 配价 %d/%d 标的, 回退接口 %d 标的", + len(closes_by_sym), + len(syms), + fallbacks, + ) + if not out_rows: + return pl.DataFrame(schema=schema) + return ( + pl.DataFrame(out_rows, schema=schema) + .unique(subset=["symbol", "trade_date"], keep="last") + .sort(["symbol", "trade_date"]) + ) + + def _load_adj_events( + self, + symset: set[str], + start_time: datetime | None, + end_time: datetime | None, + ) -> pl.DataFrame: + """事件 dump → 清洗 → 窗口过滤。返回 symbol/ex_date/dividend/bonus/allot/allot_price。 + + 清洗规则(对拍实证, 见 2026-08 验证记录): + - ex_date_ms 为上海零点戳, +8h 转日期; + - 全零事件行(疑似特殊事件, dump 未给成分)过滤; + - 同日拆行(如分红/送转各一行)按成分合并后推导, 顺序不可反; + - 配股但配股价缺失 → 无法推导, 过滤。 + """ + today = datetime.now().date() + df = self._ensure_dump(_ADJ_DUMP_KIND, "adj_factors").rename({"thscode": "symbol"}) + df = ( + df.with_columns( + pl.from_epoch(pl.col("ex_date_ms") + _SH_MS, time_unit="ms") + .dt.date() + .alias("ex_date"), + pl.col("dividend_per_share").fill_null(0.0), + pl.col("per_share_bonus").fill_null(0.0), + pl.col("allotment_ratio").fill_null(0.0), + pl.col("allotment_price").fill_null(0.0), + ) + .filter( + (pl.col("ex_date") <= today) # 未来已公告事件无前收盘, 留给滚动增量 + & ( + (pl.col("dividend_per_share") != 0) + | (pl.col("per_share_bonus") != 0) + | (pl.col("allotment_ratio") != 0) + ) + ) + .group_by("symbol", "ex_date") + .agg( + pl.col("dividend_per_share").sum().alias("dividend"), + pl.col("per_share_bonus").sum().alias("bonus"), + pl.col("allotment_ratio").sum().alias("allot"), + pl.col("allotment_price").max().alias("allot_price"), # 同日拆行共享配股价 + ) + .filter(~((pl.col("allot") > 0) & (pl.col("allot_price") <= 0))) + .filter(pl.col("symbol").is_in(sorted(symset))) + ) + if start_time is not None: + df = df.filter(pl.col("ex_date") >= start_time.date()) + if end_time is not None: + df = df.filter(pl.col("ex_date") <= end_time.date()) + return df.select("symbol", "ex_date", "dividend", "bonus", "allot", "allot_price") + + # ---- financial ---- + # 字段映射: 扶摇原始字段 → 项目 canonical 列名(TickFlow 口径, 前端财务页与回测 + # FUNDAMENTAL_FACTORS 按此消费)。映射表之外的扶摇独有字段以原名透传为扩展列。 + _INCOME_FIELD_MAP = { + "operating_income": "revenue", + "operating_costs": "operating_cost", + "sales_fee": "selling_expense", + "manage_fee": "admin_expense", + "research_and_development_expenses": "rd_expense", + "operating_profit": "operating_profit", + "interest_expenses": "financial_expense", # 近似口径: 扶摇只给利息费用 + "profit_total": "total_profit", + "income_tax_expense": "income_tax", + "net_profit": "net_income", + "parent_holder_net_profit": "net_income_attributable", + "basic_eps": "basic_eps", + } + _BALANCE_FIELD_MAP = { + "assets_total": "total_assets", + "total_current_assets": "total_current_assets", + "non_current_nets_total": "total_non_current_assets", + "cash": "cash_and_equivalents", + "accounts_receivable": "accounts_receivable", + "total_debt": "total_liabilities", + "holder_equity_total": "total_equity", + } + _CASHFLOW_FIELD_MAP = { + "act_cash_flow_net": "net_operating_cash_flow", + "invest_cash_flow_net": "net_investing_cash_flow", + "financing_cash_flow_net": "net_financing_cash_flow", + "pay_fixed_assets_etc_cash": "capex", + "cash_equivalents_net_addition": "net_cash_change", + } + # 官方指标 index_id → canonical。归母净利同比近似 tickflow net_income_yoy; + # 实测 index_id 与文档有出入(calculate_ 前缀等), 以实测为准。 + _METRICS_FIELD_MAP = { + "index_weighted_avg_roe": "roe", + "total_assets_net_ratio": "roa", + "sale_gross_margin": "gross_margin", + "sale_net_interest_ratio": "net_margin", + "assets_debt_ratio": "debt_to_asset_ratio", + "calculate_operating_income_yoy_growth_ratio": "revenue_yoy", + "calculate_parent_holder_net_profit_yoy_growth_ratio": "net_income_yoy", + "operating_cash_flow_net_divide_income": "operating_cash_to_revenue", + "inventory_turnover_ratio": "inventory_turnover", + } + + def get_financials( + self, + table: str, + symbols: list[str], + latest_only: bool = True, + ) -> pl.DataFrame: + """拉取财务数据, 映射为 canonical 列(symbol/period_end/announce_date/指标)。 + + - 三大报表: 单股单请求, latest_only 决定最近 1 期还是 8 期季报; + - metrics: 指标接口为单股单期, 恒只拉最新一期(bps 由估值快照 pb_mrq 反推, + eps_basic 顺带取自利润表); 历史各期建议切回 TickFlow 同步补齐 — + 报告期合并写入会让两源数据共存, 互不覆盖; + - shares: 扶摇无股本接口, 恒返回空(已有存量靠合并写入保留)。 + """ + if table == "shares": + logger.info("扶摇无股本接口, shares 表跳过 (已有数据保留)") + return pl.DataFrame() + if table in _STATEMENT_ENDPOINTS: + field_map = { + "income": self._INCOME_FIELD_MAP, + "balance_sheet": self._BALANCE_FIELD_MAP, + "cash_flow": self._CASHFLOW_FIELD_MAP, + }[table] + return self._financial_statements(table, field_map, symbols, latest_only) + if table == "metrics": + return self._financial_metrics(symbols) + return pl.DataFrame() + + def _financial_statements( + self, + stmt: str, + field_map: dict[str, str], + symbols: list[str], + latest_only: bool, + ) -> pl.DataFrame: + client = self._get_client() + limit = 1 if latest_only else _FINANCIAL_HISTORY_PERIODS + rows_out: list[dict] = [] + for i, sym in enumerate(symbols): + if i: + time.sleep(_HIST_INTERVAL_S) + try: + rows = client.financial_statements(stmt, sym, limit=limit) + except FuyaoError as e: + logger.warning("扶摇财务 %s %s 失败: %s", stmt, sym, e) + continue + for r in rows: + row: dict = { + "symbol": sym, + "period_end": _iso_of_ms(r.get("period_end_ms")), + "announce_date": _iso_of_ms(r.get("report_date_ms")), + } + for src, dst in field_map.items(): + row[dst] = _to_float(r.get(src)) + # 扶摇独有字段以原名透传为扩展列 (canonical 之外的增量信息) + for src, value in r.items(): + if src not in field_map and src not in row and isinstance(value, (int, float)): + row[src] = value + rows_out.append(row) + return pl.DataFrame(rows_out) if rows_out else pl.DataFrame() + + def _financial_metrics(self, symbols: list[str]) -> pl.DataFrame: + client = self._get_client() + # 指标接口按 report(yyyy-N) 单期查询 → 先用利润表 limit=1 反查每股最新披露期 + latest: dict[str, dict] = {} + for i, sym in enumerate(symbols): + if i: + time.sleep(_HIST_INTERVAL_S) + try: + rows = client.financial_statements("income", sym, limit=1) + except FuyaoError as e: + logger.warning("扶摇财务 income %s 失败: %s", sym, e) + continue + if rows: + latest[sym] = rows[0] + if not latest: + return pl.DataFrame() + bps_by_sym = self._derive_bps(sorted(latest)) + rows_out: list[dict] = [] + for sym, r in latest.items(): + quarter = _report_quarter(r.get("fiscal_period")) + report = f"{r.get('fiscal_year')}-{quarter}" if quarter else None + row: dict = { + "symbol": sym, + "period_end": _iso_of_ms(r.get("period_end_ms")), + "announce_date": _iso_of_ms(r.get("report_date_ms")), + "eps_basic": _to_float(r.get("basic_eps")), + "bps": bps_by_sym.get(sym), + } + if report: + try: + abilities = client.financial_indicators(sym, report) + except FuyaoError as e: + logger.warning("扶摇指标 %s %s 失败: %s", sym, report, e) + abilities = [] + for ability in abilities: + for ind in ability.get("indicators") or []: + index_id = ind.get("index_id") + if not index_id: + continue + value = _to_float(ind.get("value")) + if value is not None: + row[self._METRICS_FIELD_MAP.get(index_id, index_id)] = value + rows_out.append(row) + return pl.DataFrame(rows_out) if rows_out else pl.DataFrame() + + def _derive_bps(self, symbols: list[str]) -> dict[str, float]: + """估值快照 pb_mrq 与行情快照最新价同源同刻 → bps = price / pb_mrq。 + + 与财报口径 bps 可能差几个百分点(上游权益基准不完全透明), 用于补齐 + metrics.bps 使回测 pb_latest 因子可用。接口失败只影响 bps 列, 不致命。 + """ + client = self._get_client() + pb: dict[str, float] = {} + price: dict[str, float] = {} + for i in range(0, len(symbols), _VALUATION_BATCH): + if i: + time.sleep(_HIST_INTERVAL_S) + batch = symbols[i : i + _VALUATION_BATCH] + for fetch, store in ( + (client.valuations_snapshot, pb), + (client.price_snapshot_batch, price), + ): + time.sleep(_HIST_INTERVAL_S) + try: + for r in fetch(batch): + value = _to_float(r.get("pb_mrq" if store is pb else "last_price")) + code = r.get("thscode") + if code and value is not None and value != 0: + store[code] = value + except FuyaoError as e: + logger.warning("扶摇 bps 推导快照失败(%d 只): %s", len(batch), e) + return { + sym: price[sym] / pb[sym] + for sym in symbols + if sym in pb and pb[sym] and sym in price + } + + def _fetch_closes(self, symbol: str, start_d: date, end_d: date) -> dict[date, float]: + rows = self._historical_bars(symbol, start_d, end_d) + out: dict[date, float] = {} + for r in rows: + d = _date_of_ms(r.get("date_ms")) + c = _to_float(r.get("close_price")) + if d is not None and c is not None: + out[d] = c + return out + + def _closes_from_dumps( + self, bounds: pl.DataFrame + ) -> dict[str, dict[date, float]] | None: + """从本地日K dump 一次取齐全部事件标的的收盘价。 + + 每标的开窗 [first_ex-30d, last_ex](与单标的接口同窗): 10 年大 dump 为 + 主体, 10d dump 叠加补末端新鲜度(除权日当天的自检需要 ex 日收盘)。 + 返回 symbol → {date: close}; 大 dump 不可用/窗口早于其覆盖/读盘失败 + → None, 由调用方整轮回退单标的接口。 + """ + try: + lo = bounds["first_ex"].min() - timedelta(days=_PREV_CLOSE_BACKDAYS) + path = self._ensure_daily_big_dump(lo) + if path is None: + return None + sources = [self._closes_scan(pl.scan_parquet(path))] + ten = self._dump_memo.get(_DAILY10_DUMP_KIND) + if ten is not None: + sources.append(self._closes_scan(ten.lazy())) + else: + try: + p10 = self._ensure_dump_path(_DAILY10_DUMP_KIND, "daily_k_10d") + except FuyaoError as e: + logger.info("扶摇 10d dump 不可用, 配价仅用 10 年 dump: %s", e) + else: + sources.append(self._closes_scan(pl.scan_parquet(p10))) + win = bounds.select( + "symbol", + (pl.col("first_ex") - pl.duration(days=_PREV_CLOSE_BACKDAYS)).alias("lo"), + "last_ex", + ) + # concat 顺序 = 叠加优先级: 同 (symbol, date) 时 10d dump(更新鲜)覆盖大 dump + df = ( + pl.concat(sources, how="vertical_relaxed") + .join(win.lazy(), on="symbol", how="inner") + .filter( + (pl.col("date") >= pl.col("lo")) & (pl.col("date") <= pl.col("last_ex")) + ) + .unique(subset=["symbol", "date"], keep="last", maintain_order=True) + .collect() + ) + return { + f["symbol"][0]: dict( + zip(f["date"].to_list(), f["close"].to_list(), strict=True) + ) + for f in df.partition_by("symbol") + } + except Exception as e: # 缓存损坏等不致命: 回退逐标的接口 + logger.warning("扶摇除权因子本地配价失败, 回退单标的接口: %s", e) + return None + + @staticmethod + def _closes_scan(lf: pl.LazyFrame) -> pl.LazyFrame: + """dump 行 → (symbol, date, close) lazy 投影; adjusted 列存在时锁 none。""" + if "adjusted" in lf.collect_schema().names(): + lf = lf.filter(pl.col("adjusted") == "none") + return lf.select( + pl.col("thscode").alias("symbol"), + pl.from_epoch(pl.col("date_ms").cast(pl.Int64) + _SH_MS, time_unit="ms") + .dt.date() + .alias("date"), + pl.col("close_price").cast(pl.Float64).alias("close"), + ).filter(pl.col("close").is_not_null()) + # ---- 测试(设置页试拉) ---- def test_dataset(self, dataset: str, symbols: list[str] | None = None) -> dict: + if dataset in ("daily", "adj_factor"): + syms = [s for s in (symbols or [])][:3] or ["000001.SZ"] + try: + if dataset == "daily": + df = self.get_daily(syms, datetime.now() - timedelta(days=30), datetime.now()) + else: + df = self.get_adj_factors( + syms, datetime.now() - timedelta(days=365), datetime.now() + ) + except FuyaoError as e: + return {"provider": self.name, "dataset": dataset, "rows": 0, "error": str(e)} + head = df.head(5).to_dicts() + for row in head: # date/datetime → ISO 字符串, 保证 JSON 可序列化 + for k, v in list(row.items()): + if isinstance(v, (date, datetime)): + row[k] = v.isoformat() + return { + "provider": self.name, + "dataset": dataset, + "rows": df.height, + "columns": df.columns, + "preview": head, + } + if dataset == "financial": + syms = [s for s in (symbols or [])][:1] or ["600519.SH"] + try: + df = self.get_financials("metrics", syms, latest_only=True) + except FuyaoError as e: + return {"provider": self.name, "dataset": dataset, "rows": 0, "error": str(e)} + head = df.head(5).to_dicts() + return { + "provider": self.name, + "dataset": dataset, + "rows": df.height, + "columns": df.columns, + "preview": head, + } if dataset != "realtime": - return {"provider": self.name, "dataset": dataset, "rows": 0, - "error": f"扶摇插件未接入 {dataset} 数据集(自动回退 TickFlow)"} + return { + "provider": self.name, + "dataset": dataset, + "rows": 0, + "error": f"扶摇插件未接入 {dataset} 数据集(自动回退 TickFlow)", + } try: rows, count = self._get_client().snapshot_page(limit=5) except FuyaoError as e: diff --git a/backend/app/services/financial_sync.py b/backend/app/services/financial_sync.py index 7b857eb..6655496 100644 --- a/backend/app/services/financial_sync.py +++ b/backend/app/services/financial_sync.py @@ -155,6 +155,13 @@ def _sync_table( def _merge_report_history(*frames: pl.DataFrame) -> pl.DataFrame: + """按 (symbol, period_end) 合并各报告期, 同期多行逐列取最新非空值。 + + 语义(区分"覆盖"与"填空"): 每列独立取 announce_date 最新的非空值 — + 新同步行有值则覆盖旧值, 新行缺的列(如 fuyao 不提供的字段)由旧行补齐, + 实现多数据源并集共存。历史报告期不可变, 合并不会引入过期数据。 + 无 announce_date 的帧按输入顺序, 后写优先(与旧行为 keep="last" 一致)。 + """ valid = [ frame for frame in frames @@ -166,11 +173,15 @@ def _merge_report_history(*frames: pl.DataFrame) -> pl.DataFrame: pl.concat(valid, how="diagonal_relaxed") .filter(pl.col("symbol").is_not_null() & pl.col("period_end").is_not_null()) ) - # 同一 (symbol, period_end) 多条时保留 announce_date 最新一条 (业绩修正以最新公告为准)。 - if "announce_date" in merged.columns: - merged = merged.sort(["symbol", "period_end", "announce_date"], nulls_last=True) - return merged.unique(subset=["symbol", "period_end"], keep="last").sort( - ["symbol", "period_end"] + sort_keys = ["symbol", "period_end"] + ( + ["announce_date"] if "announce_date" in merged.columns else [] + ) + merged = merged.sort(sort_keys, nulls_last=True) + value_cols = [c for c in merged.columns if c not in ("symbol", "period_end")] + return ( + merged.group_by("symbol", "period_end") + .agg([pl.col(c).drop_nulls().last() for c in value_cols]) + .sort(["symbol", "period_end"]) ) diff --git a/backend/tests/test_fuyao_financial.py b/backend/tests/test_fuyao_financial.py new file mode 100644 index 0000000..e594a9c --- /dev/null +++ b/backend/tests/test_fuyao_financial.py @@ -0,0 +1,242 @@ +"""fuyao 财务适配测试 (不依赖真实网络)。 + +覆盖: 三大报表字段映射 (canonical 列名 + 扩展列透传 + ISO 日期口径)、 +latest_only 分档 (limit 1 vs 8)、metrics 组装 (eps_basic 顺带 / bps 估值反推 / +指标 index_id 映射与未知 id 透传 / 单股指标失败不弃行)、shares 恒空、 +报告期合并写入的逐列填空语义 (并集共存, 新行缺列不覆盖旧值)。 +""" + +from __future__ import annotations + +import polars as pl +import pytest + +from app.plugins.fuyao import client as fc +from app.plugins.fuyao import provider as fp +from app.plugins.fuyao.provider import FuyaoProvider +from app.services.financial_sync import _merge_report_history + + +class _FakeFinClient: + """财务端点假客户端: 记录调用入参, 按表返回预置行。""" + + def __init__( + self, + statements: dict[str, list[dict]] | None = None, + indicators: dict[str, list[dict]] | None = None, + indicator_error: Exception | None = None, + valuations: list[dict] | None = None, + prices: list[dict] | None = None, + ): + self.statements = statements or {} + self.indicators = indicators or {} + self.indicator_error = indicator_error + self.valuations = valuations or [] + self.prices = prices or [] + self.stmt_calls: list[tuple] = [] + self.ind_calls: list[str] = [] + + def financial_statements(self, stmt, thscode, limit=1): + self.stmt_calls.append((stmt, thscode, limit)) + return [dict(r, thscode=thscode) for r in self.statements.get(stmt, [])] + + def financial_indicators(self, thscode, report): + self.ind_calls.append(f"{thscode}@{report}") + if self.indicator_error: + raise self.indicator_error + return self.indicators.get(report, []) + + def valuations_snapshot(self, thscodes): + return [r for r in self.valuations if r.get("thscode") in thscodes] + + def price_snapshot_batch(self, thscodes): + return [r for r in self.prices if r.get("thscode") in thscodes] + + +def _provider_with(monkeypatch, fake): + monkeypatch.setattr( + fp, "fuyao_client", type("M", (), {"FuyaoClient": lambda **kw: fake}) + ) + monkeypatch.setattr(fp, "get_api_key", lambda: "test-key") + monkeypatch.setattr(fp, "_HIST_INTERVAL_S", 0) + return FuyaoProvider() + + +# period_end_ms: 2026-06-30 上海零点; report_date_ms: 2026-08-15 上海零点 +_INCOME_ROW = { + "period": "quarterly", + "fiscal_year": 2026, + "fiscal_period": "Q2", + "report_date_ms": 1786723200000, + "period_end_ms": 1782748800000, + "currency": "CNY", + "operating_income": 90703260964.48, + "net_profit": 46033330566.78, + "parent_holder_net_profit": 44516880421.86, + "basic_eps": 35.57, + "operating_expenses": 50000000000.0, # 扶摇独有 → 扩展列 +} + + +def test_income_mapping_canonical_columns(monkeypatch): + fake = _FakeFinClient(statements={"income": [_INCOME_ROW]}) + provider = _provider_with(monkeypatch, fake) + df = provider.get_financials("income", ["600519.SH"], latest_only=True) + row = df.to_dicts()[0] + assert row["symbol"] == "600519.SH" + assert row["period_end"] == "2026-06-30" + assert row["announce_date"] == "2026-08-15" + assert row["revenue"] == pytest.approx(_INCOME_ROW["operating_income"]) + assert row["net_income"] == pytest.approx(_INCOME_ROW["net_profit"]) + assert row["net_income_attributable"] == pytest.approx( + _INCOME_ROW["parent_holder_net_profit"] + ) + assert row["basic_eps"] == 35.57 + # 扩展列: 扶摇独有数字字段原名透传; 字符串元数据不透传 + assert row["operating_expenses"] == 50000000000.0 + assert "thscode" not in df.columns and "ticker" not in df.columns + # 原始名不残留 (已映射字段) + assert "operating_income" not in df.columns and "net_profit" not in df.columns + + +def test_statements_limit_latest_vs_history(monkeypatch): + fake = _FakeFinClient(statements={"income": [_INCOME_ROW]}) + provider = _provider_with(monkeypatch, fake) + provider.get_financials("income", ["600519.SH"], latest_only=True) + assert fake.stmt_calls == [("income", "600519.SH", 1)] + provider.get_financials("income", ["600519.SH"], latest_only=False) + assert fake.stmt_calls[-1] == ("income", "600519.SH", fp._FINANCIAL_HISTORY_PERIODS) + + +def test_balance_and_cashflow_mapping(monkeypatch): + balance = { + "period_end_ms": 1782748800000, + "report_date_ms": 1786723200000, + "assets_total": 309050784569.31, + "total_debt": 46954432394.95, + "holder_equity_total": 262096352174.36, + } + cashflow = { + "period_end_ms": 1782748800000, + "report_date_ms": 1786723200000, + "act_cash_flow_net": 92000000000.0, + "invest_cash_flow_net": -3000000000.0, + "pay_dividends_profits_interest_cash": 64000000000.0, # 扶摇独有 → 扩展列 + } + fake = _FakeFinClient( + statements={"balance_sheet": [balance], "cash_flow": [cashflow]} + ) + provider = _provider_with(monkeypatch, fake) + bal = provider.get_financials("balance_sheet", ["600519.SH"]).to_dicts()[0] + assert bal["total_assets"] == pytest.approx(balance["assets_total"]) + assert bal["total_liabilities"] == pytest.approx(balance["total_debt"]) + assert bal["total_equity"] == pytest.approx(balance["holder_equity_total"]) + cf = provider.get_financials("cash_flow", ["600519.SH"]).to_dicts()[0] + assert cf["net_operating_cash_flow"] == pytest.approx(92000000000.0) + assert cf["net_investing_cash_flow"] == pytest.approx(-3000000000.0) + assert cf["pay_dividends_profits_interest_cash"] == pytest.approx(64000000000.0) + + +_METRICS_ABILITIES = [ + { + "ability": "profitability", + "indicators": [ + {"index_id": "index_weighted_avg_roe", "value": "16.7500"}, + {"index_id": "sale_gross_margin", "value": "89.5552"}, + ], + }, + { + "ability": "growth", + # 实测 index_id 与文档有出入; 未映射 id 原名透传 + "indicators": [{"index_id": "fixed_asset_invest_expansion_ratio", "value": "2.12587300"}], + }, + { + "ability": "solvency", + "indicators": [{"index_id": "earned_interest_multiple", "value": None}], # null → 不写列 + }, +] + + +def test_metrics_assembly(monkeypatch): + fake = _FakeFinClient( + statements={"income": [_INCOME_ROW]}, + indicators={"2026-2": _METRICS_ABILITIES}, + valuations=[{"thscode": "600519.SH", "pb_mrq": 6.455055}], + prices=[{"thscode": "600519.SH", "last_price": 1297.4}], + ) + provider = _provider_with(monkeypatch, fake) + df = provider.get_financials("metrics", ["600519.SH"], latest_only=True) + row = df.to_dicts()[0] + assert row["period_end"] == "2026-06-30" + assert row["announce_date"] == "2026-08-15" + assert row["eps_basic"] == 35.57 # 顺带取自利润表 + assert row["bps"] == pytest.approx(1297.4 / 6.455055) # 估值反推 + assert row["roe"] == pytest.approx(16.75) # 字符串 → float + assert row["gross_margin"] == pytest.approx(89.5552) + assert row["fixed_asset_invest_expansion_ratio"] == pytest.approx(2.125873) + assert "earned_interest_multiple" not in df.columns # 全空指标不成列 + assert fake.ind_calls == ["600519.SH@2026-2"] # report 由利润表最新期反推 + + +def test_metrics_indicator_failure_keeps_row(monkeypatch): + """指标端点单股失败 (如未披露期 code=5003) → 行仍写入 (eps/bps 保留)。""" + fake = _FakeFinClient( + statements={"income": [_INCOME_ROW]}, + indicator_error=fc.FuyaoError("code=5003"), + valuations=[{"thscode": "600519.SH", "pb_mrq": 6.455055}], + prices=[{"thscode": "600519.SH", "last_price": 1297.4}], + ) + provider = _provider_with(monkeypatch, fake) + df = provider.get_financials("metrics", ["600519.SH"], latest_only=True) + row = df.to_dicts()[0] + assert row["symbol"] == "600519.SH" + assert row["eps_basic"] == 35.57 + assert "roe" not in df.columns + + +def test_metrics_skips_symbol_without_income(monkeypatch): + fake = _FakeFinClient(statements={"income": []}) + provider = _provider_with(monkeypatch, fake) + assert provider.get_financials("metrics", ["600519.SH"]).is_empty() + + +def test_shares_returns_empty(monkeypatch): + provider = _provider_with(monkeypatch, _FakeFinClient()) + assert provider.get_financials("shares", ["600519.SH"]).is_empty() + + +def test_merge_fills_missing_cells_from_old_rows(): + """逐列填空: 同报告期新行缺的列由旧行补齐, 有值则覆盖 (并集共存语义)。""" + old = pl.DataFrame({ + "symbol": ["600519.SH", "600519.SH"], + "period_end": ["2026-03-31", "2026-06-30"], + "announce_date": ["2026-04-20", "2026-08-10"], + "diluted_eps": [68.1, 70.2], # tickflow 提供, fuyao 没有 + "net_income": [280.0, 460.0], + }) + new = pl.DataFrame({ + "symbol": ["600519.SH"], + "period_end": ["2026-06-30"], + "announce_date": ["2026-08-15"], # 更晚公告 → 该期以新行为基准 + "net_income": [461.5], # 修正值覆盖 + # diluted_eps 缺失 → 由旧行 70.2 补齐 + "bps": [200.99], # fuyao 扩展列, 旧行没有 + }) + merged = _merge_report_history(old, new).to_dicts() + assert len(merged) == 2 + q2 = next(r for r in merged if r["period_end"] == "2026-06-30") + assert q2["net_income"] == pytest.approx(461.5) # 新值覆盖 + assert q2["diluted_eps"] == pytest.approx(70.2) # 旧行补齐 + assert q2["bps"] == pytest.approx(200.99) # 扩展列并入 + q1 = next(r for r in merged if r["period_end"] == "2026-03-31") + assert q1["diluted_eps"] == pytest.approx(68.1) # 未触碰期原样保留 + # 旧公告覆盖新公告的倒序场景: announce 早的行不覆盖晚的 + reversed_new = pl.DataFrame({ + "symbol": ["600519.SH"], + "period_end": ["2026-06-30"], + "announce_date": ["2026-08-01"], + "net_income": [999.0], + }) + q2b = _merge_report_history(old, reversed_new).to_dicts()[1] + # 公告更晚的 old 行 (08-10) 胜出, 早公告的新行不覆盖 → 业绩修正以最新公告为准 + assert q2b["net_income"] == pytest.approx(460.0) diff --git a/backend/tests/test_fuyao_provider.py b/backend/tests/test_fuyao_provider.py index 7e158d0..c5d3693 100644 --- a/backend/tests/test_fuyao_provider.py +++ b/backend/tests/test_fuyao_provider.py @@ -1,11 +1,18 @@ """FuyaoProvider 契约与单位标准化测试。 -不依赖真实网络: 用假 FuyaoClient 返回样例快照页, 验证字段映射、 -百分数→小数制转换 (CONTRIBUTING §3.1)、分页合并、软失败、 -能力声明 (未声明数据集回退 tickflow) 与设置页试拉。 +不依赖真实网络: 用假 FuyaoClient 返回样例快照页/历史K线/dump, 验证字段映射、 +单位口径 (CONTRIBUTING §3.1: 百分数→小数、volume 股→手)、日K取数分档 +(10d dump / 单标的接口)、除权因子推导 (交易所公式 + half-up 舍入 + +同日合并 + 涨跌停自检)、能力声明 (未声明数据集回退 tickflow) 与设置页试拉。 """ + from __future__ import annotations +import calendar +import itertools +from datetime import date, datetime, timedelta + +import polars as pl import pytest from app.plugins.fuyao import client as fc @@ -16,7 +23,13 @@ from app.plugins.fuyao.provider import FuyaoProvider class _FakeClient: """按调用次数返回预置页, 记录调用供分页断言。snapshot_all 同真实客户端语义。""" - def __init__(self, pages: list[list[dict]], count: int, error: Exception | None = None, server_ts: int = 0): + def __init__( + self, + pages: list[list[dict]], + count: int, + error: Exception | None = None, + server_ts: int = 0, + ): self.pages = pages self.count = count self.error = error @@ -73,7 +86,9 @@ def _row(thscode: str = "600519.SH", **over): def _provider_with(monkeypatch, pages, count=None, error=None, **fake_kwargs): - fake = _FakeClient(pages, count if count is not None else sum(len(p) for p in pages), error, **fake_kwargs) + fake = _FakeClient( + pages, count if count is not None else sum(len(p) for p in pages), error, **fake_kwargs + ) monkeypatch.setattr(fp, "fuyao_client", type("M", (), {"FuyaoClient": lambda **kw: fake})) monkeypatch.setattr(fp, "get_api_key", lambda: "test-key") return FuyaoProvider(), fake @@ -81,6 +96,7 @@ def _provider_with(monkeypatch, pages, count=None, error=None, **fake_kwargs): # ---- 单位与字段映射 ---- + def test_snapshot_units_and_field_mapping(monkeypatch): """核心口径: price_change_ratio_pct 百分数 → change_pct 小数制 (1.72 → 0.0172)。""" provider, _ = _provider_with(monkeypatch, [[_row()]]) @@ -94,7 +110,8 @@ def test_snapshot_units_and_field_mapping(monkeypatch): assert r["open"] == 1460.0 assert r["high"] == 1490.5 assert r["low"] == 1455.0 - assert r["volume"] == 1234500 + # volume 单位股(1234500) → 内部契约手: floor(1234500/100) = 12345 + assert r["volume"] == 12345 assert r["amount"] == 1.83e9 assert r["timestamp"] > 0 # 快照不提供的字段必须为 None, 不启发式伪造 @@ -130,10 +147,12 @@ def test_all_rows_unrecognized_returns_empty_with_no_fake_data(monkeypatch): # ---- 客户端信封解析 (实测结构 vs 文档示例) ---- + def _patch_http(monkeypatch, payload, status_code=200): class _Resp: def json(self): return payload + _Resp.status_code = status_code class _Http: @@ -148,11 +167,14 @@ def _patch_http(monkeypatch, payload, status_code=200): def test_client_parses_real_world_envelope(monkeypatch): """实测信封(2026-08): data={timestamp, total, item}。""" - _patch_http(monkeypatch, { - "code": 0, "message": "success", - "data": {"timestamp": 1787542612000, "total": 2, - "item": [_row(), _row("000001.SZ")]}, - }) + _patch_http( + monkeypatch, + { + "code": 0, + "message": "success", + "data": {"timestamp": 1787542612000, "total": 2, "item": [_row(), _row("000001.SZ")]}, + }, + ) c = fc.FuyaoClient(api_key="k") rows, total = c.snapshot_page() assert total == 2 and len(rows) == 2 @@ -161,10 +183,14 @@ def test_client_parses_real_world_envelope(monkeypatch): def test_client_parses_documented_envelope(monkeypatch): """官方文档示例信封: data={count, data}。""" - _patch_http(monkeypatch, { - "code": 0, "message": "OK", - "data": {"count": 3, "data": [_row()]}, - }) + _patch_http( + monkeypatch, + { + "code": 0, + "message": "OK", + "data": {"count": 3, "data": [_row()]}, + }, + ) c = fc.FuyaoClient(api_key="k") rows, total = c.snapshot_page() assert total == 3 and len(rows) == 1 @@ -179,13 +205,20 @@ def test_client_raises_on_error_code(monkeypatch): # ---- 字段名兼容与服务端时间戳 ---- + def test_doc_style_field_names_fallback(monkeypatch): """文档示例字段名 (highest_price/lowest_price/prev_close_price) 也能映射。""" row = { - "thscode": "600519.SH", "last_price": 1480.0, "price_change": 25.0, - "price_change_ratio_pct": 1.72, "open_price": 1460.0, - "highest_price": 1490.5, "lowest_price": 1455.0, "prev_close_price": 1455.0, - "volume": 1234500, "turnover": 1.83e9, + "thscode": "600519.SH", + "last_price": 1480.0, + "price_change": 25.0, + "price_change_ratio_pct": 1.72, + "open_price": 1460.0, + "highest_price": 1490.5, + "lowest_price": 1455.0, + "prev_close_price": 1455.0, + "volume": 1234500, + "turnover": 1.83e9, } provider, _ = _provider_with(monkeypatch, [[row]]) r = provider.get_realtime()[0] @@ -206,6 +239,7 @@ def test_realtime_falls_back_to_local_time_without_server_ts(monkeypatch): # ---- 分页 ---- + def test_snapshot_pagination_merges_pages(monkeypatch): page1 = [_row(f"{600000 + i}.SH") for i in range(2)] page2 = [_row(f"{688000 + i}.SH") for i in range(1)] @@ -224,8 +258,11 @@ def test_snapshot_stops_when_page_empty(monkeypatch): # ---- 软失败 ---- + def test_realtime_error_returns_empty_list(monkeypatch): - provider, _ = _provider_with(monkeypatch, [[]], error=fc.FuyaoError("扶摇接口错误 code=4001: 频率超限")) + provider, _ = _provider_with( + monkeypatch, [[]], error=fc.FuyaoError("扶摇接口错误 code=4001: 频率超限") + ) assert provider.get_realtime() == [] @@ -236,19 +273,23 @@ def test_client_requires_api_key(): # ---- 能力声明与注册 ---- -def test_datasets_declaration_realtime_only(): - """只声明 realtime; 其他数据集 provider_has_dataset 必须为 False (回退 tickflow)。""" + +def test_datasets_declaration(): + """声明 realtime/daily/adj_factor/financial; minute 未声明 (回退 tickflow)。""" config = FuyaoProvider().config assert "realtime" in config.datasets - assert "daily" not in config.datasets + assert "daily" in config.datasets + assert "adj_factor" in config.datasets + assert "financial" in config.datasets assert "minute" not in config.datasets - assert "financial" not in config.datasets # ---- API Key 解析 (secrets.json > .env, 对齐 tickflow 语义) ---- + def test_get_api_key_secrets_store_takes_priority(monkeypatch): from app import secrets_store + monkeypatch.delenv(fp.API_KEY_ENV, raising=False) monkeypatch.setattr(secrets_store, "load", lambda: {fp.SECRETS_FIELD: "sk-from-ui"}) assert fp.get_api_key() == "sk-from-ui" @@ -256,6 +297,7 @@ def test_get_api_key_secrets_store_takes_priority(monkeypatch): def test_get_api_key_falls_back_to_env(monkeypatch): from app import secrets_store + monkeypatch.setenv(fp.API_KEY_ENV, "sk-from-env") monkeypatch.setattr(secrets_store, "load", lambda: {}) assert fp.get_api_key() == "sk-from-env" @@ -263,6 +305,7 @@ def test_get_api_key_falls_back_to_env(monkeypatch): def test_availability_accepts_secrets_store_key(monkeypatch): from app import secrets_store + monkeypatch.delenv(fp.API_KEY_ENV, raising=False) monkeypatch.setattr(secrets_store, "load", lambda: {fp.SECRETS_FIELD: "sk-from-ui"}) assert fp.availability() == (True, "ok") @@ -270,6 +313,7 @@ def test_availability_accepts_secrets_store_key(monkeypatch): def test_availability_requires_env_key(monkeypatch): from app import secrets_store + monkeypatch.delenv(fp.API_KEY_ENV, raising=False) monkeypatch.setattr(secrets_store, "load", lambda: {}) ok, reason = fp.availability() @@ -278,6 +322,7 @@ def test_availability_requires_env_key(monkeypatch): # ---- 先探后存 (probe_api_key) ---- + def _patch_client_cls(monkeypatch, fake): monkeypatch.setattr(fp, "fuyao_client", type("M", (), {"FuyaoClient": lambda **kw: fake})) @@ -289,7 +334,10 @@ def test_probe_api_key_ok(monkeypatch): def test_probe_api_key_invalid_key(monkeypatch): - _patch_client_cls(monkeypatch, _FakeClient([[]], 0, error=fc.FuyaoError("扶摇接口错误 code=1001: 无效 api key"))) + _patch_client_cls( + monkeypatch, + _FakeClient([[]], 0, error=fc.FuyaoError("扶摇接口错误 code=1001: 无效 api key")), + ) ok, reason = fp.probe_api_key("sk-bad") assert ok is False and "无效" in reason @@ -298,13 +346,16 @@ def test_loader_probe_plugin_key_dispatch(monkeypatch): import app.plugins.fuyao.provider as provider_mod from app.data_providers.custom import loader - monkeypatch.setattr(provider_mod, "probe_api_key", lambda key: (True, "ok") if key == "good" else (False, "bad")) + monkeypatch.setattr( + provider_mod, "probe_api_key", lambda key: (True, "ok") if key == "good" else (False, "bad") + ) assert loader.probe_plugin_key("fuyao", "good") == (True, "ok") assert loader.probe_plugin_key("fuyao", "bad") == (False, "bad") def test_loader_probe_plugin_key_unsupported_plugin(): from app.data_providers.custom import loader + # stock-sdk 未声明 api_key_env → 不支持界面配 Key ok, reason = loader.probe_plugin_key("stocksdk", "x") assert ok is False and "不支持" in reason @@ -314,13 +365,16 @@ def test_loader_probe_plugin_key_unsupported_plugin(): # ---- 保存/清除端点 (直接调用 handler, 先探后存语义) ---- + def test_save_plugin_key_invalid_key_not_persisted(monkeypatch): from app.api import settings as settings_api from app.data_providers import custom as custom_sources saved: dict = {} monkeypatch.setattr(custom_sources, "probe_plugin_key", lambda n, k: (False, "Key 无效")) - monkeypatch.setattr(settings_api.secrets_store, "save", lambda updates: saved.update(updates) or updates) + monkeypatch.setattr( + settings_api.secrets_store, "save", lambda updates: saved.update(updates) or updates + ) out = settings_api.save_plugin_key(settings_api.PluginKeyIn(plugin="fuyao", api_key="bad")) assert out["ok"] is False and out["reason"] == "invalid" assert saved == {} # 无效 Key 不落盘 @@ -333,10 +387,16 @@ def test_save_plugin_key_valid_persists_and_rescans(monkeypatch): saved: dict = {} reloaded = [] monkeypatch.setattr(custom_sources, "probe_plugin_key", lambda n, k: (True, "ok")) - monkeypatch.setattr(settings_api.secrets_store, "save", lambda updates: saved.update(updates) or updates) - monkeypatch.setattr(settings_api.secrets_store, "mask", lambda key, prefix=4, suffix=4: "abcd••••wxyz") + monkeypatch.setattr( + settings_api.secrets_store, "save", lambda updates: saved.update(updates) or updates + ) + monkeypatch.setattr( + settings_api.secrets_store, "mask", lambda key, prefix=4, suffix=4: "abcd••••wxyz" + ) monkeypatch.setattr(custom_sources, "load_all", lambda: reloaded.append(1)) - monkeypatch.setattr(custom_sources, "list_plugins", lambda: [{"name": "fuyao", "available": True}]) + monkeypatch.setattr( + custom_sources, "list_plugins", lambda: [{"name": "fuyao", "available": True}] + ) out = settings_api.save_plugin_key(settings_api.PluginKeyIn(plugin="fuyao", api_key="good-key")) assert out["ok"] is True assert saved == {"fuyao_api_key": "good-key"} # 字段名与 provider.SECRETS_FIELD 一致 @@ -352,18 +412,21 @@ def test_clear_plugin_key(monkeypatch): monkeypatch.setattr(custom_sources, "is_builtin", lambda n: n == "fuyao") monkeypatch.setattr(settings_api.secrets_store, "clear", lambda *keys: cleared.extend(keys)) monkeypatch.setattr(custom_sources, "load_all", lambda: None) - monkeypatch.setattr(custom_sources, "list_plugins", lambda: [{"name": "fuyao", "available": False}]) + monkeypatch.setattr( + custom_sources, "list_plugins", lambda: [{"name": "fuyao", "available": False}] + ) out = settings_api.clear_plugin_key("fuyao") assert out["ok"] is True and out["plugin_available"] is False assert cleared == ["fuyao_api_key"] -def test_manifest_declares_realtime_dataset(): +def test_manifest_declares_datasets(): from app.data_providers.custom import loader + manifest = loader.plugin_manifest("fuyao") assert manifest is not None assert manifest["entry"] == "app.plugins.fuyao.provider:FuyaoProvider" - assert "realtime" in (manifest.get("datasets") or []) + assert {"realtime", "daily", "adj_factor"} <= set(manifest.get("datasets") or []) assert manifest.get("runtime") == "none" assert manifest.get("api_key_env") == fp.API_KEY_ENV @@ -374,6 +437,7 @@ def test_hidden_plugin_not_registered(): 用合成清单验证: hidden: true 的插件不注册、不在数据源页展示。 """ from app.data_providers.custom import loader + manifest = loader.plugin_manifest("fuyao") assert manifest is not None and not manifest.get("hidden"), ( "fuyao 应保持可见; 如需重新隐藏请在 plugin.yaml 声明 hidden 并更新本测试" @@ -389,6 +453,7 @@ def test_hidden_plugin_not_registered(): # ---- 插件 Key 脱敏展示 (与 TickFlow Key 契约一致) ---- + def test_plugin_key_masked_from_secrets_then_env(monkeypatch): """api_key_masked 随插件状态返回: secrets.json 优先, .env 兜底, 未配置为空。 @@ -419,6 +484,7 @@ def test_plugin_key_masked_from_secrets_then_env(monkeypatch): # ---- 设置页试拉 ---- + def test_test_dataset_realtime_preview(monkeypatch): provider, _ = _provider_with(monkeypatch, [[_row(), _row("000001.SZ")]], count=5400) out = provider.test_dataset("realtime") @@ -438,3 +504,851 @@ def test_close_is_idempotent(monkeypatch): provider, _ = _provider_with(monkeypatch, [[_row()]]) provider.close() provider.close() + + +# ===================================================================== +# daily / adj_factor (2026-08 接入) +# ===================================================================== + + +def _sh_ms(d: date) -> int: + """交易日 → 扶摇口径 ms(该日上海零点), 与 provider._ms_of_date 同式。""" + return (calendar.timegm(d.timetuple()) - 28_800) * 1000 + + +def _bar(d: date, close: float, volume: float = 1_612_611, open_: float | None = None): + """historical/dump 原始行(价格元, volume 股)。""" + return { + "date_ms": _sh_ms(d), + "open_price": open_ if open_ is not None else round(close * 0.99, 2), + "high_price": round(close * 1.01, 2), + "low_price": round(close * 0.98, 2), + "close_price": close, + "volume": volume, + "turnover": close * volume, + } + + +def _dump_bar(sym: str, d: date, close: float) -> dict: + """daily dump 行(11 列形状, 含 thscode/adjusted)。""" + return { + "thscode": sym, + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(d, close), + } + + +class _FakeHistClient: + """单标的历史K线假客户端: 按窗口过滤预置 bars, 记录调用参数。""" + + def __init__(self, bars_by_symbol: dict | None = None, error_syms: tuple = ()): + self.bars = bars_by_symbol or {} + self.errors = set(error_syms) + self.calls: list[dict] = [] + + def historical_kline(self, thscode, start_ms, end_ms, adjust="none"): + self.calls.append({"thscode": thscode, "start": start_ms, "end": end_ms, "adjust": adjust}) + if thscode in self.errors: + raise fc.FuyaoError("扶摇接口错误 code=4001: 频率超限") + return [b for b in self.bars.get(thscode, []) if start_ms <= b["date_ms"] <= end_ms] + + def dump_download_url(self, dump_kind): + raise fc.FuyaoError("测试环境无 dump") + + def close(self): + pass + + +def _hist_provider(monkeypatch, fake: _FakeHistClient, allow_dumps: bool = False) -> FuyaoProvider: + p = FuyaoProvider() + monkeypatch.setattr(fp, "fuyao_client", type("M", (), {"FuyaoClient": lambda **kw: fake})) + monkeypatch.setattr(fp, "get_api_key", lambda: "test-key") + monkeypatch.setattr(fp, "_HIST_INTERVAL_S", 0.0) + if not allow_dumps: + # 默认禁用 dump 档(单标的接口路径测试用); 大 dump 测试传 allow_dumps=True + p._ensure_daily_big_dump = lambda start_d: None # type: ignore[assignment] + return p + + +def _daily10_dump(rows: list[dict]) -> pl.DataFrame: + """11 列 daily-k-10d dump 形状。""" + return pl.DataFrame( + rows, + schema={ + "thscode": pl.String, + "currency": pl.String, + "interval": pl.String, + "adjusted": pl.String, + "date_ms": pl.Int64, + "open_price": pl.Float64, + "high_price": pl.Float64, + "low_price": pl.Float64, + "close_price": pl.Float64, + "volume": pl.Float64, + "turnover": pl.Float64, + }, + ) + + +def _adj_dump(events: list[tuple]) -> pl.DataFrame: + """8 列 adjustment-factors dump 形状。events: (thscode, ex_date, D, S, AR, AP)。""" + rows = [ + { + "thscode": s, + "ticker": s.split(".")[0], + "ex_date_ms": _sh_ms(d), + "dividend_per_share": dv, + "per_share_bonus": bn, + "allotment_ratio": ar, + "allotment_price": ap, + "currency": "CNY", + } + for s, d, dv, bn, ar, ap in events + ] + return pl.DataFrame( + rows, + schema={ + "thscode": pl.String, + "ticker": pl.String, + "ex_date_ms": pl.Int64, + "dividend_per_share": pl.Float64, + "per_share_bonus": pl.Float64, + "allotment_ratio": pl.Float64, + "allotment_price": pl.Float64, + "currency": pl.String, + }, + ) + + +def _adj_provider(monkeypatch, events: list[tuple], bars_by_symbol: dict) -> FuyaoProvider: + p = _hist_provider(monkeypatch, _FakeHistClient(bars_by_symbol)) + p._dump_memo[fp._ADJ_DUMP_KIND] = _adj_dump(events) + return p + + +# ---- daily: 单标的接口路径 ---- + + +def test_daily_api_units_volume_shares_to_lots(monkeypatch): + """核心口径: 原始价 + volume 股→手 floor(/100) + 上海零点时区。""" + bars = [ + _bar(date(2026, 8, 26), 11.10), + _bar(date(2026, 8, 27), 11.05), + _bar(date(2026, 8, 28), 11.65), + ] + provider = _hist_provider(monkeypatch, _FakeHistClient({"000001.SZ": bars})) + # 窗口跨度 28 天 > _RECENT_DUMP_DAYS → 走单标的接口 + df = provider.get_daily(["000001.SZ"], datetime(2026, 8, 1), datetime(2026, 8, 28)) + assert df.height == 3 + assert df.columns == ["symbol", "date", "open", "high", "low", "close", "volume", "amount"] + assert df.schema["date"] == pl.Date + assert df["date"].to_list() == [date(2026, 8, 26), date(2026, 8, 27), date(2026, 8, 28)] + assert df["volume"].to_list() == [16126.0] * 3 # 1,612,611 股 → 16126 手 + assert df["close"].to_list() == [11.10, 11.05, 11.65] + assert df["amount"].to_list() == [b["turnover"] for b in bars] + + +def test_daily_api_adjust_locked_to_none(monkeypatch): + """adjust=none 锁定: 服务端默认 forward, 官方前复权序列不可用(事件间逐日漂移)。""" + provider = _hist_provider( + monkeypatch, _FakeHistClient({"000001.SZ": [_bar(date(2026, 8, 27), 11.05)]}) + ) + provider.get_daily(["000001.SZ"], datetime(2026, 8, 1), datetime(2026, 8, 28)) + fake = provider._get_client() + assert fake.calls and all(c["adjust"] == "none" for c in fake.calls) + + +def test_daily_api_deep_window_splits_into_10y_chunks(monkeypatch): + """超 10 年窗口自动分片(25.6 年 → 3 片), 分片连续不重叠且各 ≤10 年。""" + provider = _hist_provider(monkeypatch, _FakeHistClient({"000001.SZ": []})) + provider.get_daily(["000001.SZ"], datetime(2001, 1, 1), datetime(2026, 8, 28)) + calls = provider._get_client().calls + assert len(calls) == 3 + for c in calls: + assert c["end"] - c["start"] < fp._HIST_MAX_SPAN_MS + for a, b in itertools.pairwise(calls): + assert b["start"] == a["end"] + 1 # 连续且不重叠 + + +def test_daily_api_soft_fail_per_symbol(monkeypatch): + """单标的失败软跳过, 不阻断其余标的。""" + bars = {"000001.SZ": [_bar(date(2026, 8, 27), 11.05)]} + provider = _hist_provider(monkeypatch, _FakeHistClient(bars, error_syms=("600519.SH",))) + df = provider.get_daily(["600519.SH", "000001.SZ"], datetime(2026, 8, 1), datetime(2026, 8, 28)) + assert df["symbol"].unique().to_list() == ["000001.SZ"] + + +def test_daily_empty_symbols_or_non_stock_returns_empty(monkeypatch): + provider = _hist_provider(monkeypatch, _FakeHistClient({})) + assert provider.get_daily([], None, None).is_empty() + assert provider.get_daily( + ["510300.SH"], datetime(2026, 8, 1), datetime(2026, 8, 28), asset_type="etf" + ).is_empty() + + +# ---- daily: 10d dump 路径 ---- + + +def _recent_daily_dump() -> pl.DataFrame: + rows = [] + for d in [date(2026, 8, 25), date(2026, 8, 26), date(2026, 8, 27), date(2026, 8, 28)]: + b = _bar(d, 11.0 + hash(d) % 3, volume=97_570_170) + rows.append( + {"thscode": "000001.SZ", "currency": "CNY", "interval": "1d", "adjusted": "none", **b} + ) + rows.append( + { + "thscode": "600519.SH", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(d, 1480.0, volume=1_234_500), + } + ) + # 停牌行(open=high=0)应被过滤 + rows.append( + { + "thscode": "000001.SZ", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2026, 8, 27), 11.0, open_=0.0), + } + ) + rows[-1]["high_price"] = 0.0 + return _daily10_dump(rows) + + +def test_daily_recent_window_uses_dump_without_api_calls(monkeypatch): + """近端窗口: 一次 dump 覆盖全部标的, 不打单标的接口; volume 股→手; 停牌行过滤。""" + provider = _hist_provider(monkeypatch, _FakeHistClient({})) + provider._dump_memo[fp._DAILY10_DUMP_KIND] = _recent_daily_dump() + progress = [] + df = provider.get_daily( + ["000001.SZ"], + datetime(2026, 8, 26), + datetime(2026, 8, 28), + on_chunk_done=lambda c, t: progress.append((c, t)), + ) + assert provider._get_client().calls == [] # 未走单标的接口 + assert progress == [(1, 1)] + assert set(df["symbol"].to_list()) == {"000001.SZ"} # 未请求的标的不混入 + assert df["volume"].max() == 975_701.0 # 97,570,170 股 → 975,701 手 + assert date(2026, 8, 26) in df["date"].to_list() + # 停牌行(open=high=0)被 filter_halt_days 剔除, 08-27 仅保留正常行 + assert df.filter(pl.col("date") == date(2026, 8, 27)).height == 1 + + +def test_daily_dump_stale_falls_back_to_api(monkeypatch): + """dump 末端落后且 end 是工作日 → 视为有缺口, 回退单标的接口。""" + bars = {"000001.SZ": [_bar(date(2026, 8, 27), 11.05)]} + provider = _hist_provider(monkeypatch, _FakeHistClient(bars)) + provider._dump_memo[fp._DAILY10_DUMP_KIND] = _recent_daily_dump() # max=08-28 + provider.get_daily(["000001.SZ"], datetime(2026, 8, 24), datetime(2026, 9, 4)) # 09-04 周五 + assert provider._get_client().calls != [] + + +def test_daily_dump_weekend_end_covered(monkeypatch): + """end 为周末且紧随 dump 末端 → 自然缺口, 仍走 dump。""" + provider = _hist_provider(monkeypatch, _FakeHistClient({})) + provider._dump_memo[fp._DAILY10_DUMP_KIND] = _recent_daily_dump() # max=周五 08-28 + df = provider.get_daily(["000001.SZ"], datetime(2026, 8, 26), datetime(2026, 8, 30)) # 周日 + assert provider._get_client().calls == [] + assert not df.is_empty() + + +def test_daily_dump_rejects_non_raw_adjustment(monkeypatch): + """防御: dump 变为复权口径(adjusted != none)时拒绝输出, 不污染原始K线库。""" + dump = _recent_daily_dump().with_columns(pl.lit("forward").alias("adjusted")) + provider = _hist_provider(monkeypatch, _FakeHistClient({})) + provider._dump_memo[fp._DAILY10_DUMP_KIND] = dump + assert provider.get_daily( + ["000001.SZ"], datetime(2026, 8, 26), datetime(2026, 8, 28) + ).is_empty() + + +def test_daily_old_window_bypasses_dump(monkeypatch): + """深窗口且 dump 不可用 → 单标的接口兜底。""" + provider = _hist_provider(monkeypatch, _FakeHistClient({"000001.SZ": []})) + df = provider.get_daily(["000001.SZ"], datetime(2020, 1, 1), datetime(2026, 8, 28)) + assert df.is_empty() + assert provider._get_client().calls # 走了单标的接口 + + +# ---- daily: 10 年全量 dump 档 ---- + + +def _bigdump_provider( + monkeypatch, tmp_path, big_rows: list[dict], ten_memo: pl.DataFrame | None = None +): + """缓存目录重定向到 tmp_path 并放置 daily_k__.parquet; 单标的接口若被调用会被断言暴露。""" + provider = _hist_provider(monkeypatch, _FakeHistClient({}), allow_dumps=True) + monkeypatch.setattr(fp, "_cache_dir", lambda: tmp_path) + _daily10_dump(big_rows).write_parquet(tmp_path / "daily_k__20260101.parquet") + if ten_memo is not None: + provider._dump_memo[fp._DAILY10_DUMP_KIND] = ten_memo + return provider + + +def test_daily_deep_window_uses_big_dump(monkeypatch, tmp_path): + """深窗口走 10 年 dump: 不打单标的接口, volume 股→手, 未请求标的不混入。""" + rows = [ + { + "thscode": "000001.SZ", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2026, 1, 5), 11.0, volume=97_570_170), + }, + { + "thscode": "000001.SZ", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2026, 6, 12), 11.7, volume=1_234_500), + }, + { + "thscode": "000001.SZ", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2026, 8, 14), 11.4, volume=1_234_500), + }, + { + "thscode": "600519.SH", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2026, 8, 14), 1480.0, volume=1_234_500), + }, + ] + provider = _bigdump_provider(monkeypatch, tmp_path, rows) + progress = [] + df = provider.get_daily( + ["000001.SZ"], + datetime(2026, 1, 5), + datetime(2026, 8, 14), + on_chunk_done=lambda c, t: progress.append((c, t)), + ) + assert provider._get_client().calls == [] + assert progress == [(1, 1)] + assert set(df["symbol"].to_list()) == {"000001.SZ"} + assert df["date"].to_list() == [date(2026, 1, 5), date(2026, 6, 12), date(2026, 8, 14)] + assert df["volume"].to_list() == [975_701.0, 12_345.0, 12_345.0] + + +def test_daily_big_dump_tail_filled_by_10d(monkeypatch, tmp_path): + """大 dump 末端缺口(dmax 旧)由 10d dump 补尾, 两段拼接无缝。""" + big_rows = [ + { + "thscode": "000001.SZ", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2026, 8, 10), 10.9), + }, + { + "thscode": "000001.SZ", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2026, 8, 14), 11.0), + }, + ] + ten_rows = [ + { + "thscode": "000001.SZ", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2026, 8, 15), 11.2), + }, + { + "thscode": "000001.SZ", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2026, 8, 28), 11.65), + }, + ] + provider = _bigdump_provider(monkeypatch, tmp_path, big_rows, ten_memo=_daily10_dump(ten_rows)) + df = provider.get_daily(["000001.SZ"], datetime(2026, 8, 10), datetime(2026, 8, 28)) + assert df["date"].to_list() == [ + date(2026, 8, 10), + date(2026, 8, 14), + date(2026, 8, 15), + date(2026, 8, 28), + ] + assert df["close"].to_list() == [10.9, 11.0, 11.2, 11.65] + + +def test_daily_big_dump_midgap_falls_back(monkeypatch, tmp_path): + """大 dump 与 10d dump 之间有中段缺口 → 整体回退, 不拼缺口数据。""" + big_rows = [ + { + "thscode": "000001.SZ", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2026, 8, 1), 10.9), + }, + { + "thscode": "000001.SZ", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2026, 8, 5), 11.0), + }, + ] + # 10d 从 08-27 起, 与大 dump 末端 08-05 之间有中段缺口 + ten_rows = [ + { + "thscode": "000001.SZ", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2026, 8, 27), 11.5), + }, + ] + provider = _bigdump_provider(monkeypatch, tmp_path, big_rows, ten_memo=_daily10_dump(ten_rows)) + df = provider.get_daily(["000001.SZ"], datetime(2026, 8, 1), datetime(2026, 8, 28)) + assert df.is_empty() # 兜底走单标的接口(fake 为空) → 不返回缺口拼接数据 + + +def test_daily_big_dump_start_before_dmin_falls_back(monkeypatch, tmp_path): + """窗口早于大 dump 起点(>10 年) → 单标的接口兜底。""" + rows = [ + { + "thscode": "000001.SZ", + "currency": "CNY", + "interval": "1d", + "adjusted": "none", + **_bar(date(2020, 1, 2), 10.0), + }, + ] + provider = _bigdump_provider(monkeypatch, tmp_path, rows) + # 缓存不覆盖起点 → 会尝试拉最新 release; 模拟"最新 release 就是这份缓存文件" + provider._ensure_dump_path = lambda kind, prefix: tmp_path / "daily_k__20260101.parquet" # type: ignore[assignment] + bars = {"000001.SZ": [_bar(date(2015, 6, 1), 9.0)]} + provider._dump_memo.clear() + monkeypatch.setattr( + fp, "fuyao_client", type("M", (), {"FuyaoClient": lambda **kw: _FakeHistClient(bars)}) + ) + df = provider.get_daily(["000001.SZ"], datetime(2015, 1, 1), datetime(2015, 12, 31)) + assert df.height == 1 and df["close"][0] == 9.0 # 单标的兜底成功 + + +def test_tail_ok_weekend_tolerance(): + """末端容忍: 周末/节假日的自然缺口(≤3 天且 end 是周末)不算缺失。""" + assert fp._tail_ok(date(2026, 8, 28), date(2026, 8, 28)) is True + assert fp._tail_ok(date(2026, 8, 29), date(2026, 8, 28)) is True # 周六 + assert fp._tail_ok(date(2026, 8, 30), date(2026, 8, 28)) is True # 周日 + assert fp._tail_ok(date(2026, 8, 31), date(2026, 8, 28)) is False # 周一(可能有行情) + assert fp._tail_ok(date(2026, 9, 5), date(2026, 8, 28)) is False # 隔了一周 + + +# ---- adj_factor: 推导 ---- + + +def test_adj_dividend_factor_derivation(monkeypatch): + """纯分红: P=32.8, D=0.68 → ref=32.12 → factor=32.8/32.12 (与 tickflow 对拍值一致)。""" + events = [("600519.SH", date(2026, 6, 12), 0.68, 0.0, 0.0, 0.0)] + bars = {"600519.SH": [_bar(date(2026, 6, 11), 32.8), _bar(date(2026, 6, 12), 32.0)]} + provider = _adj_provider(monkeypatch, events, bars) + df = provider.get_adj_factors(["600519.SH"], datetime(2026, 1, 1), datetime(2026, 12, 31)) + assert df.height == 1 + row = df.row(0, named=True) + assert row["trade_date"] == date(2026, 6, 12) + assert row["ex_factor"] == pytest.approx(32.8 / 32.12) + + +def test_adj_bonus_and_allotment_formula(monkeypatch): + """送转+配股混合: P=5.23, AR=0.4, AP=3.36 → ref=(5.23+1.344)/1.4=4.6957→4.70(half-up)。""" + events = [("300176.SZ", date(2026, 8, 21), 0.0, 0.0, 0.4, 3.36)] + bars = {"300176.SZ": [_bar(date(2026, 8, 20), 5.23), _bar(date(2026, 8, 21), 4.68)]} + provider = _adj_provider(monkeypatch, events, bars) + df = provider.get_adj_factors(["300176.SZ"], datetime(2026, 1, 1), datetime(2026, 12, 31)) + assert df.height == 1 + assert df["ex_factor"][0] == pytest.approx(5.23 / 4.70) + + +def test_adj_half_up_rounding_not_bankers(monkeypatch): + """舍入口径: x=2.625 → half-up 2.63 (银行家舍入会给 2.62, 对拍实证偏离)。""" + events = [("600000.SH", date(2026, 6, 12), 8.0, 0.0, 0.0, 0.0)] + bars = {"600000.SH": [_bar(date(2026, 6, 11), 10.625), _bar(date(2026, 6, 12), 2.70)]} + provider = _adj_provider(monkeypatch, events, bars) + df = provider.get_adj_factors(["600000.SH"], datetime(2026, 1, 1), datetime(2026, 12, 31)) + assert df["ex_factor"][0] == pytest.approx(10.625 / 2.63) + + +def test_adj_price_limit_selfcheck_drops_impossible_factor(monkeypatch): + """涨跌停自检: 复权后除权日涨幅超板限的事件剔除, 不让坏因子落库。""" + events = [("600000.SH", date(2026, 6, 12), 8.0, 0.0, 0.0, 0.0)] + bars = {"600000.SH": [_bar(date(2026, 6, 11), 10.625), _bar(date(2026, 6, 12), 4.05)]} # +54% + provider = _adj_provider(monkeypatch, events, bars) + df = provider.get_adj_factors(["600000.SH"], datetime(2026, 1, 1), datetime(2026, 12, 31)) + assert df.is_empty() + + +def test_adj_merges_same_day_components_before_derivation(monkeypatch): + """同日拆行(分红行+送股行)必须先合并再推导。""" + events = [ + ("000812.SZ", date(2026, 6, 12), 0.3, 0.0, 0.0, 0.0), + ("000812.SZ", date(2026, 6, 12), 0.0, 0.1, 0.0, 0.0), + ] + bars = {"000812.SZ": [_bar(date(2026, 6, 11), 10.0), _bar(date(2026, 6, 12), 9.0)]} + provider = _adj_provider(monkeypatch, events, bars) + df = provider.get_adj_factors(["000812.SZ"], datetime(2026, 1, 1), datetime(2026, 12, 31)) + assert df.height == 1 # 合并成一个事件 + assert df["ex_factor"][0] == pytest.approx(10.0 / 8.82) # ref=(10-0.3)/1.1=8.8181→8.82 + + +def test_adj_filters_zero_rows_and_missing_allot_price(monkeypatch): + """全零行与配股价缺失行过滤, 不参与推导。""" + events = [ + ("000001.SZ", date(2026, 6, 12), 0.0, 0.0, 0.0, 0.0), # 全零 + ("000002.SZ", date(2026, 6, 12), 0.0, 0.0, 0.3, 0.0), # 配股无价 + ("600519.SH", date(2026, 6, 12), 0.68, 0.0, 0.0, 0.0), # 正常 + ] + bars = { + "600519.SH": [_bar(date(2026, 6, 11), 32.8), _bar(date(2026, 6, 12), 32.0)], + "000001.SZ": [_bar(date(2026, 6, 11), 10.0), _bar(date(2026, 6, 12), 10.0)], + "000002.SZ": [_bar(date(2026, 6, 11), 8.0), _bar(date(2026, 6, 12), 8.0)], + } + provider = _adj_provider(monkeypatch, events, bars) + df = provider.get_adj_factors(["000001.SZ", "000002.SZ", "600519.SH"], None, None) + assert df["symbol"].unique().to_list() == ["600519.SH"] + + +def test_adj_skips_future_announced_events(monkeypatch): + """未来已公告事件无前收盘, 跳过留给滚动增量窗口。""" + future = date.today() + timedelta(days=5) + events = [("600519.SH", future, 0.68, 0.0, 0.0, 0.0)] + provider = _adj_provider(monkeypatch, events, {"600519.SH": []}) + df = provider.get_adj_factors(["600519.SH"], datetime(2026, 1, 1), None) + assert df.is_empty() + + +def test_adj_window_filter(monkeypatch): + events = [ + ("600519.SH", date(2026, 1, 10), 0.5, 0.0, 0.0, 0.0), + ("600519.SH", date(2026, 6, 12), 0.68, 0.0, 0.0, 0.0), + ] + bars = { + "600519.SH": [ + _bar(date(2026, 1, 9), 30.0), + _bar(date(2026, 1, 10), 30.0), + _bar(date(2026, 6, 11), 32.8), + _bar(date(2026, 6, 12), 32.0), + ] + } + provider = _adj_provider(monkeypatch, events, bars) + df = provider.get_adj_factors(["600519.SH"], datetime(2026, 3, 1), datetime(2026, 12, 31)) + assert df["trade_date"].to_list() == [date(2026, 6, 12)] + + +def test_adj_output_schema_sorted_and_deduped(monkeypatch): + events = [ + ("000001.SZ", date(2026, 6, 12), 0.36, 0.0, 0.0, 0.0), + ("600519.SH", date(2026, 6, 12), 0.68, 0.0, 0.0, 0.0), + ("600519.SH", date(2025, 6, 20), 0.5, 0.0, 0.0, 0.0), + ] + bars = { + "600519.SH": [ + _bar(date(2025, 6, 19), 31.0), + _bar(date(2025, 6, 20), 31.0), + _bar(date(2026, 6, 11), 32.8), + _bar(date(2026, 6, 12), 32.0), + ], + "000001.SZ": [_bar(date(2026, 6, 11), 11.7), _bar(date(2026, 6, 12), 11.5)], + } + provider = _adj_provider(monkeypatch, events, bars) + df = provider.get_adj_factors(["000001.SZ", "600519.SH"], None, None) + assert df.columns == ["symbol", "trade_date", "ex_factor"] + assert df.schema["trade_date"] == pl.Date and df.schema["ex_factor"] == pl.Float64 + assert df.height == 3 + assert df.sort(["symbol", "trade_date"]).equals(df) # 已排序 + assert df["ex_factor"].min() > 1.0 # 分红因子必 >1 + + +def test_adj_etf_and_empty_symbols_return_empty(monkeypatch): + provider = _adj_provider(monkeypatch, [], {}) + empty = provider.get_adj_factors([], None, None) + assert empty.is_empty() and empty.columns == ["symbol", "trade_date", "ex_factor"] + etf = provider.get_adj_factors(["510300.SH"], None, None, asset_type="etf") + assert etf.is_empty() + + +def test_adj_dump_unavailable_returns_empty(monkeypatch): + """dump 加载失败软返回空(不阻断管道), 不抛异常。""" + provider = _hist_provider(monkeypatch, _FakeHistClient({})) + + def _boom(dump_kind, cache_prefix): + raise fc.FuyaoError("dump 下载网络失败") + + provider._ensure_dump = _boom # type: ignore[assignment] + df = provider.get_adj_factors(["600519.SH"], None, None) + assert df.is_empty() and df.columns == ["symbol", "trade_date", "ex_factor"] + + +def test_adj_progress_callback_per_symbol(monkeypatch): + events = [ + ("000001.SZ", date(2026, 6, 12), 0.36, 0.0, 0.0, 0.0), + ("600519.SH", date(2026, 6, 12), 0.68, 0.0, 0.0, 0.0), + ] + bars = { + "600519.SH": [_bar(date(2026, 6, 11), 32.8), _bar(date(2026, 6, 12), 32.0)], + "000001.SZ": [_bar(date(2026, 6, 11), 11.7), _bar(date(2026, 6, 12), 11.5)], + } + provider = _adj_provider(monkeypatch, events, bars) + progress = [] + provider.get_adj_factors( + ["600519.SH", "000001.SZ"], None, None, on_chunk_done=lambda c, t: progress.append((c, t)) + ) + assert progress == [(1, 2), (2, 2)] + + +# ---- adj_factor: 本地 dump 配价(2026-08 优化) ---- + + +def test_adj_closes_from_local_dump_no_http(monkeypatch, tmp_path): + """前收盘从本地日K dump 一次取齐: 零单标的请求, 因子值与接口路径同公式。""" + events = [("000001.SZ", date(2026, 6, 12), 0.3, 0.0, 0.0, 0.0)] + big_rows = [ + _dump_bar("000001.SZ", date(2026, 5, 10), 10.2), # 覆盖窗口起点(first_ex-30d) + _dump_bar("000001.SZ", date(2026, 6, 11), 10.625), + _dump_bar("000001.SZ", date(2026, 6, 12), 9.6), + ] + provider = _bigdump_provider(monkeypatch, tmp_path, big_rows) + provider._dump_memo[fp._ADJ_DUMP_KIND] = _adj_dump(events) + progress = [] + df = provider.get_adj_factors( + ["000001.SZ"], + datetime(2026, 1, 1), + datetime(2026, 12, 31), + on_chunk_done=lambda c, t: progress.append((c, t)), + ) + assert provider._get_client().calls == [] # 零 HTTP + assert progress == [(1, 1)] + assert df.height == 1 + assert df["trade_date"][0] == date(2026, 6, 12) + assert df["ex_factor"][0] == pytest.approx(10.625 / fp._ref_price(10.625, 0.3, 0.0, 0.0, 0.0)) + + +def test_adj_missing_symbol_in_dump_falls_back_to_http(monkeypatch, tmp_path): + """dump 缺价的标的回退单标的接口; dump 内标的仍零请求。""" + events = [ + ("000001.SZ", date(2026, 6, 12), 0.3, 0.0, 0.0, 0.0), + ("600519.SH", date(2026, 6, 12), 1.0, 0.0, 0.0, 0.0), + ] + big_rows = [ + _dump_bar("000001.SZ", date(2026, 5, 10), 10.2), + _dump_bar("000001.SZ", date(2026, 6, 11), 10.625), + _dump_bar("000001.SZ", date(2026, 6, 12), 9.6), + ] + provider = _bigdump_provider(monkeypatch, tmp_path, big_rows) + provider._dump_memo[fp._ADJ_DUMP_KIND] = _adj_dump(events) + provider._get_client().bars["600519.SH"] = [ + _bar(date(2026, 6, 11), 1500.0), + _bar(date(2026, 6, 12), 1400.0), + ] + df = provider.get_adj_factors( + ["000001.SZ", "600519.SH"], datetime(2026, 1, 1), datetime(2026, 12, 31) + ) + assert [c["thscode"] for c in provider._get_client().calls] == ["600519.SH"] + assert set(df["symbol"].to_list()) == {"000001.SZ", "600519.SH"} + + +def test_adj_ten_day_dump_overlays_fresh_ex_close(monkeypatch, tmp_path): + """大 dump 缺除权日收盘(发布滞后)时由 10d dump 叠加补, 涨跌停自检仍生效。""" + events = [("000001.SZ", date(2026, 8, 28), 8.0, 0.0, 0.0, 0.0)] # 异常大额分红 + big_rows = [ + _dump_bar("000001.SZ", date(2026, 7, 25), 10.0), # 覆盖窗口起点 + _dump_bar("000001.SZ", date(2026, 8, 27), 10.625), # 大 dump 只到除权前一日 + ] + ten_rows = [_dump_bar("000001.SZ", date(2026, 8, 28), 4.05)] # +54% 超涨跌停 + provider = _bigdump_provider(monkeypatch, tmp_path, big_rows, ten_memo=_daily10_dump(ten_rows)) + provider._dump_memo[fp._ADJ_DUMP_KIND] = _adj_dump(events) + df = provider.get_adj_factors(["000001.SZ"], datetime(2026, 1, 1), datetime(2026, 12, 31)) + assert df.is_empty() # 自检剔除(无叠加时 ex 日收盘缺失, 该行会被保留) + + +def test_adj_local_pricing_corrupt_dump_falls_back(monkeypatch, tmp_path): + """本地配价读盘异常不致命: 整轮回退单标的接口。""" + events = [("000001.SZ", date(2026, 6, 12), 0.3, 0.0, 0.0, 0.0)] + big_rows = [ + _dump_bar("000001.SZ", date(2026, 5, 10), 10.2), + _dump_bar("000001.SZ", date(2026, 6, 11), 10.625), + _dump_bar("000001.SZ", date(2026, 6, 12), 9.6), + ] + provider = _bigdump_provider(monkeypatch, tmp_path, big_rows) + provider._dump_memo[fp._ADJ_DUMP_KIND] = _adj_dump(events) + + def _boom(_lf): # 签名对齐被替换的 _closes_scan + raise OSError("disk error") + + provider._closes_scan = _boom # type: ignore[assignment] + provider._get_client().bars["000001.SZ"] = [ + _bar(date(2026, 6, 11), 10.625), + _bar(date(2026, 6, 12), 9.6), + ] + df = provider.get_adj_factors(["000001.SZ"], datetime(2026, 1, 1), datetime(2026, 12, 31)) + assert df.height == 1 + assert provider._get_client().calls # 回退了接口 + + +# ---- 客户端: 历史日K / dump ---- + + +class _RecordingHttp: + def __init__(self, payload): + self.payload = payload + self.calls: list[tuple] = [] + + def get(self, path, params=None): + self.calls.append((path, params)) + resp = type("R", (), {"json": lambda self_: self.payload, "status_code": 200})() + return resp + + def close(self): + pass + + +def test_client_historical_kline_sends_params(monkeypatch): + http = _RecordingHttp({"code": 0, "data": {"item": [_bar(date(2026, 8, 28), 11.65)]}}) + monkeypatch.setattr(fc.httpx, "Client", lambda **kw: http) + c = fc.FuyaoClient(api_key="k") + rows = c.historical_kline("000001.SZ", 1, 2, adjust="none") + assert len(rows) == 1 + path, params = http.calls[0] + assert path == "/api/a-share/prices/historical" + assert params == { + "thscode": "000001.SZ", + "interval": "1d", + "adjust": "none", + "start": 1, + "end": 2, + } + + +def test_client_historical_kline_missing_item(monkeypatch): + http = _RecordingHttp({"code": 0, "data": {}}) + monkeypatch.setattr(fc.httpx, "Client", lambda **kw: http) + c = fc.FuyaoClient(api_key="k") + assert c.historical_kline("000001.SZ", 1, 2) == [] + + +def test_client_dump_download_url(monkeypatch): + info = { + "presigned_url": "https://o.thsi.cn/.../releases/20260828/x.parquet?sig=1", + "presigned_url_expires_at": "2026-08-29T12:00:00+08:00", + "expires_in_seconds": 300, + } + http = _RecordingHttp({"code": 0, "data": info}) + monkeypatch.setattr(fc.httpx, "Client", lambda **kw: http) + c = fc.FuyaoClient(api_key="k") + assert c.dump_download_url("adjustment-factors") == info + assert http.calls[0][0] == "/api/dump/market-dumps/adjustment-factors/download-url" + + +def test_client_download_dump_atomic_and_key_never_leaked(monkeypatch, tmp_path): + """/S3 预签名下载不带 X-api-key; 成功后原子改名, .part 清理。""" + http = _RecordingHttp( + { + "code": 0, + "data": { + "presigned_url": "https://o.thsi.cn/fuyao/market-dump/adj/releases/20260828/a.parquet?sig=1" + }, + } + ) + monkeypatch.setattr(fc.httpx, "Client", lambda **kw: http) + + captured: dict = {} + + class _Resp: + status_code = 200 + + def iter_bytes(self, n): + yield b"parquet-bytes-" + + class _Ctx: + def __enter__(self): + return _Resp() + + def __exit__(self, *a): + return False + + def _stream(method, url, **kw): + captured["method"], captured["url"], captured["kw"] = method, url, kw + return _Ctx() + + monkeypatch.setattr(fc.httpx, "stream", _stream) + dest = tmp_path / "adj__20260828.parquet" + fc.FuyaoClient(api_key="k").download_dump("adjustment-factors", dest) + assert captured["method"] == "GET" and "o.thsi.cn" in captured["url"] + assert "headers" not in captured["kw"] or "X-api-key" not in ( + captured["kw"].get("headers") or {} + ) + assert dest.read_bytes() == b"parquet-bytes-" + assert not list(tmp_path.glob("*.part")) + + +def test_client_download_dump_http_error(monkeypatch, tmp_path): + http = _RecordingHttp({"code": 0, "data": {"presigned_url": "https://o.thsi.cn/x.parquet"}}) + monkeypatch.setattr(fc.httpx, "Client", lambda **kw: http) + + class _Resp: + status_code = 403 + + class _Ctx: + def __enter__(self): + return _Resp() + + def __exit__(self, *a): + return False + + monkeypatch.setattr(fc.httpx, "stream", lambda m, u, **kw: _Ctx()) + dest = tmp_path / "x.parquet" + with pytest.raises(fc.FuyaoError, match="403"): + fc.FuyaoClient(api_key="k").download_dump("adjustment-factors", dest) + assert not dest.exists() and not list(tmp_path.glob("*.part")) + + +# ---- 时区与 release 解析 ---- + + +def test_shanghai_midnight_ms_roundtrip(): + """上海零点戳 ↔ 日期: +8h 换算不得偏移一天。""" + d = date(2026, 6, 12) + assert fp._date_of_ms(_sh_ms(d)) == d + assert fp._ms_of_date(d) == _sh_ms(d) + + +def test_release_of_extracts_from_presigned_url(): + url = "https://o.thsi.cn/x/market-dump/adj/releases/20260828/a.parquet?sig=1" + assert fp._release_of(url) == "20260828" + assert fp._release_of("https://no-release-here/x") == "unknown" + + +# ---- 设置页试拉: daily / adj_factor ---- + + +def test_test_dataset_daily_preview(monkeypatch): + bars = {"000001.SZ": [_bar(date(2026, 8, 27), 11.05), _bar(date(2026, 8, 28), 11.65)]} + provider = _hist_provider(monkeypatch, _FakeHistClient(bars)) + out = provider.test_dataset("daily", ["000001.SZ"]) + assert out["provider"] == "fuyao" and out["dataset"] == "daily" + assert out["rows"] == 2 + assert out["preview"][0]["date"] == "2026-08-27" # date → ISO 字符串 + + +def test_test_dataset_adj_factor_preview(monkeypatch): + events = [("600519.SH", date(2026, 6, 12), 0.68, 0.0, 0.0, 0.0)] + bars = {"600519.SH": [_bar(date(2026, 6, 11), 32.8), _bar(date(2026, 6, 12), 32.0)]} + provider = _adj_provider(monkeypatch, events, bars) + out = provider.test_dataset("adj_factor", ["600519.SH"]) + assert out["rows"] == 1 + assert out["preview"][0]["trade_date"] == "2026-06-12" + assert out["preview"][0]["ex_factor"] == pytest.approx(32.8 / 32.12) diff --git a/docs/custom-data-source.md b/docs/custom-data-source.md index 954e5fe..104d701 100644 --- a/docs/custom-data-source.md +++ b/docs/custom-data-source.md @@ -296,7 +296,9 @@ cp docs/examples/custom-data-source/mock_source.yaml data/data_sources/mock_sour 分钟K (minute): symbol = 股票代码 - datetime = 时间戳 (YYYY-MM-DD HH:MM:SS) + # datetime 必须是北京时间墙钟 (如 2026-08-28 09:35:00), 不要返回 UTC; + # 入口守卫会自动纠偏 UTC 特征帧, 但契约仍要求源头写对 + datetime = 北京时间墙钟 (YYYY-MM-DD HH:MM:SS) open / high / low / close = OHLC volume = 成交量 amount = 成交额 diff --git a/docs/plugin-development.md b/docs/plugin-development.md index 55f0743..9e1ccc9 100644 --- a/docs/plugin-development.md +++ b/docs/plugin-development.md @@ -147,7 +147,7 @@ class MyProvider: def get_minute(self, symbols, start_time, end_time, asset_type="stock", on_chunk_done=None, freq="1m") -> pl.DataFrame: - """分钟K: [symbol, datetime, open, high, low, close, volume, amount]""" + """分钟K: [symbol, datetime(北京墙钟), open, high, low, close, volume, amount]""" def get_realtime(self) -> list[dict]: """全市场实时快照 → list[dict]。失败软返回 [], 不抛异常(不阻断轮询线程)。""" @@ -164,6 +164,21 @@ class MyProvider: 返回 error 字段说明会回退 TickFlow。""" ``` +### get_minute 的 datetime 时区契约 + +`datetime` 必须是**北京时间墙钟**(naive,如 `2026-08-28 09:35:00`),与日K的 +`date` 语义对齐;不要返回 UTC 或带时区的时间。前端分时图按交易时段时轴 +(09:30–11:30 / 13:00–15:00)映射每根K线,UTC 口径的帧会导致全部点位落在时轴外、 +分时图空白。 + +入口守卫(`kline_sync._enforce_minute_beijing_wallclock`)对所有分钟源强制归一: +带时区 → 自动换算成北京墙钟;naive 但整体呈 UTC 特征(如 01:30)→ 自动 +8 纠偏并 +记日志;完全无法识别的口径 → 拒收并回退 TickFlow。契约仍要求源头写对,守卫只是兜底。 + +可选类属性 `minute_history_days = 5` 声明 1 分钟历史深度(交易日);未声明视为 +深历史(TickFlow 基准)。浅源(如 stock-sdk 免费分时仅保留最近 5 个交易日)声明后, +个股分时档位自动收窄为可行选项并默认 5 日,深源默认 20 日。 + ### 异常语义 | 方法 | 失败行为 | @@ -213,7 +228,7 @@ class MyConfig: 插件 PR 必须带契约测试(CONTRIBUTING §9), **不依赖真实网络与 API Key**——用假 Client/桥接注入。以 `backend/tests/test_fuyao_provider.py` 为范本, 至少覆盖: -1. 字段映射与单位转换: 百分数→小数制、缺失字段按口径推导、缺失字段置 None 不伪造 +1. 字段映射与单位转换: 百分数→小数制、volume 股→手、*ms 零点戳时区换算、缺失字段按口径推导、缺失字段置 None 不伪造 2. 接口响应结构变体: 实测结构 vs 官方文档示例双兼容(供应商文档与实际不一致是常态) 3. 分页: 多页合并、空页终止、页数上限 4. 软失败: 接口报错返回 []; 整页 schema 变化有告警而非静默空数据 @@ -229,10 +244,10 @@ uv run --extra dev python -m ruff check app/plugins// tests/test_