mirror of
https://ghfast.top/https://github.com/aeroxw/tick-stock-panel.git
synced 2026-09-12 14:24:15 +08:00
merge: feat/full-minute-dataset (全量分钟数据集开放) 合入 main
分支已含 main 全部提交, 合并前验证: 涉改 106 用例 + 前端 tsc 零错误 + 全量 1340 passed (test_mining_manager 偶发为存量, 合并前 main 可复现)。
This commit is contained in:
+2
-2
@@ -122,7 +122,7 @@
|
||||
数据源已经插件化。任何通用功能都必须通过 provider 能力和标准化数据集访问数据,不能把 TickFlow SDK 调用硬编码到策略、监控、回测、API 或前端流程中。
|
||||
|
||||
- 使用现有的 `get_provider()`、`provider_has_dataset()` 和 preferences 路由能力。
|
||||
- 支持的数据集包括但不限于 `daily`、`adj_factor`、`minute`、`realtime`、`financial`;新增数据集应先定义清晰的输入输出契约。
|
||||
- 支持的数据集包括但不限于 `daily`、`adj_factor`、`minute`、`full_minute`(盘中全市场分钟落盘)、`realtime`、`financial`;新增数据集应先定义清晰的输入输出契约。
|
||||
- provider 负责把供应商字段、单位、日期和代码格式转换为内部标准格式。
|
||||
- 上层服务依赖标准字段和能力声明,不依赖供应商响应结构。
|
||||
- 只有明确标注为 TickFlow 专属的功能才可以直接依赖 TickFlow,并且不得影响其他 provider。
|
||||
@@ -138,7 +138,7 @@
|
||||
- 各页面能力门控统一以矩阵的 `usable` 为准(生效源当前能否真正提供该能力),不是 TickFlow 套餐视角;缺能力提示统一引导到数据源配置。
|
||||
- 能力层中立:通用界面(侧栏徽章、能力路由卡、各页门控提示)不得出现 TickFlow 档位/订阅词汇;档位信息只在 TickFlow 专属详情卡展示。provider 名称作为路由事实可以出现。
|
||||
- 每个能力独立路由,禁止跟随/派生特殊值(`same_as_daily` 已下线);存量非法偏好值由 preferences getter 回退默认自愈,不做迁移。
|
||||
- 边界注记:分时监控由分钟能力兜底(`intraday_monitor_support`),不单设分时能力;`depth5` 已进矩阵但插件数据集白名单暂未开放,当前仅 TickFlow 提供。
|
||||
- 边界注记:分时监控由分钟能力兜底(`intraday_monitor_support`),不单设分时能力;`full_minute`(全量分钟)数据集已开放插件/自定义源声明;`depth5` 已进矩阵但插件数据集白名单暂未开放,当前仅 TickFlow 提供。
|
||||
- 实时指数为产品级固定契约,不走路由矩阵:展示层(侧栏指数条、市场总览)固定核心四只(`backend/app/services/index_const.py` 单一权威:上证/深成/创业板/科创综指),后端各消费方与前端 Layout 引用同一份定义不建副本;指数页保留但标的固定为核心四只(无全指数搜索/浏览,`/api/index/list`、`/api/index/search` 已下线);侧栏指数多选配置已下线,相关偏好(`realtime_index_symbols`/`sidebar_index_symbols`/`indices_nav_pinned`/`realtime_pull_index`/`realtime_index_mode`)已删除。监控规则的指数标的不受限——quote_service 把核心四只 + 启用规则的指数并入显式拉取。
|
||||
- 自定义源指数补充协议:A 股快照普遍不含指数(fuyao 实测无指数,指数在其独立端点)。provider 可实现可选方法 `get_realtime_indices(symbols) -> list[dict]`(record 结构与 realtime 一致),quote_service 在自定义源分支鸭子类型调用补拉;未实现的源指数缓存为空,由本地日K兜底接管。fuyao 指数快照有连坐语义——请求混入未知代码整批失败,插件侧必须先行过滤不支持的后缀(如 `.BJ`)。
|
||||
|
||||
|
||||
@@ -412,6 +412,7 @@ class DataProvidersIn(BaseModel):
|
||||
daily_data_provider: str | None = None
|
||||
adj_factor_provider: str | None = None
|
||||
minute_data_provider: str | None = None
|
||||
full_minute_data_provider: str | None = None
|
||||
depth5_data_provider: str | None = None
|
||||
realtime_data_provider: str | None = None
|
||||
financial_data_provider: str | None = None
|
||||
@@ -506,6 +507,7 @@ def get_preferences() -> dict:
|
||||
"daily_data_provider": preferences.get_daily_data_provider(),
|
||||
"adj_factor_provider": preferences.get_adj_factor_provider(),
|
||||
"minute_data_provider": preferences.get_minute_data_provider(),
|
||||
"full_minute_data_provider": preferences.get_full_minute_data_provider(),
|
||||
"minute_history_days": _minute_history_days(),
|
||||
"depth5_data_provider": preferences.get_depth5_data_provider(),
|
||||
"realtime_data_provider": preferences.get_realtime_data_provider(),
|
||||
@@ -786,6 +788,7 @@ def update_data_providers(req: DataProvidersIn, request: Request) -> dict:
|
||||
"daily_data_provider": preferences.get_daily_data_provider(),
|
||||
"adj_factor_provider": preferences.get_adj_factor_provider(),
|
||||
"minute_data_provider": preferences.get_minute_data_provider(),
|
||||
"full_minute_data_provider": preferences.get_full_minute_data_provider(),
|
||||
"depth5_data_provider": preferences.get_depth5_data_provider(),
|
||||
"realtime_data_provider": preferences.get_realtime_data_provider(),
|
||||
"financial_data_provider": preferences.get_financial_provider(),
|
||||
|
||||
@@ -81,11 +81,11 @@ CAPABILITY_REGISTRY: list[dict] = [
|
||||
"id": "full_minute",
|
||||
"label": "全量分钟",
|
||||
"desc": "盘中全市场当日分钟落盘 (冷启动全天 + 标的池增量)",
|
||||
"field": None,
|
||||
"field": "full_minute_data_provider",
|
||||
"default": "tickflow",
|
||||
"tf_tier": "expert",
|
||||
# intraday.universe 能力 (TickFlow Expert 专有): 插件契约不开放此数据集,
|
||||
# 生效源恒为 TickFlow — 不可路由, 无对应 provider 偏好字段
|
||||
# TickFlow 侧需 Expert 档; 插件/自定义源声明 full_minute 数据集即可提供
|
||||
# (插件实现 get_intraday_batch / 可选 get_intraday_latest, YAML 仅修复轮)
|
||||
},
|
||||
]
|
||||
|
||||
|
||||
@@ -423,7 +423,7 @@ def _sanitize_for_yaml(config: dict) -> dict:
|
||||
|
||||
datasets_out: dict = {}
|
||||
for ds_name, ds_cfg in (config.get("datasets") or {}).items():
|
||||
if ds_name not in {"daily", "adj_factor", "realtime", "minute", "financial"}:
|
||||
if ds_name not in {"daily", "adj_factor", "realtime", "minute", "full_minute", "financial"}:
|
||||
continue
|
||||
if not isinstance(ds_cfg, dict):
|
||||
continue
|
||||
|
||||
@@ -30,6 +30,8 @@ _REQUIRED = {
|
||||
"adj_factor": {"symbol", "trade_date", "ex_factor"},
|
||||
"realtime": {"symbol", "last_price", "prev_close", "open", "high", "low", "volume"},
|
||||
"minute": {"symbol", "datetime", "open", "high", "low", "close", "volume", "amount"},
|
||||
# full_minute (全量分钟) 与 minute 同形: 当日窗口批量拉取, 字段映射一致
|
||||
"full_minute": {"symbol", "datetime", "open", "high", "low", "close", "volume", "amount"},
|
||||
# financial 字段由数据源决定, 只要求能映射出 symbol
|
||||
"financial": {"symbol"},
|
||||
}
|
||||
@@ -116,7 +118,7 @@ class GenericHTTPProvider:
|
||||
errors.append(f"{dataset}: pct_unit 必须是 percent 或 decimal")
|
||||
if dataset != "realtime":
|
||||
request_params = [cfg.symbols_param, cfg.start_param, cfg.end_param]
|
||||
if dataset == "minute":
|
||||
if dataset in {"minute", "full_minute"}:
|
||||
request_params.extend(
|
||||
name for name in (cfg.asset_type_param, cfg.freq_param) if name
|
||||
)
|
||||
@@ -205,7 +207,38 @@ class GenericHTTPProvider:
|
||||
配置的参数名注入请求 (GET → params, POST → body), 用于上游需区分
|
||||
stock/ETF/index 或固定频率的场景。
|
||||
"""
|
||||
cfg = self._dataset("minute")
|
||||
return self._fetch_minute_dataset(
|
||||
"minute", symbols, start_time, end_time, asset_type, freq, on_chunk_done,
|
||||
)
|
||||
|
||||
def get_intraday_batch(
|
||||
self,
|
||||
symbols: list[str],
|
||||
count: int = 300, # noqa: ARG002 — 与插件契约对齐, YAML 源按时间窗口取全天
|
||||
asset_type: AssetType = "stock",
|
||||
) -> pl.DataFrame:
|
||||
"""全量分钟修复轮: 按当日窗口批量拉取 full_minute 数据集 (chunked + rpm 限速)。
|
||||
|
||||
与 get_minute 同形 (字段映射/归一一致), 区别仅在数据集名与窗口由调用方
|
||||
传当日值。稳态增量 (get_intraday_latest) YAML 声明式源不提供 — 服务自动
|
||||
降级为仅修复轮模式并放慢节奏。
|
||||
"""
|
||||
start = datetime.now().replace(hour=0, minute=0, second=0, microsecond=0)
|
||||
return self._fetch_minute_dataset(
|
||||
"full_minute", symbols, start, datetime.now(), asset_type, "1m", None,
|
||||
)
|
||||
|
||||
def _fetch_minute_dataset(
|
||||
self,
|
||||
ds_name: str,
|
||||
symbols: list[str],
|
||||
start_time: datetime | None,
|
||||
end_time: datetime | None,
|
||||
asset_type: AssetType = "stock",
|
||||
freq: str = "1m",
|
||||
on_chunk_done: Callable[[int, int], None] | None = None,
|
||||
) -> pl.DataFrame:
|
||||
cfg = self._dataset(ds_name)
|
||||
override: dict[str, Any] = {}
|
||||
if cfg.asset_type_param:
|
||||
override[cfg.asset_type_param] = asset_type
|
||||
@@ -280,7 +313,7 @@ class GenericHTTPProvider:
|
||||
start_time = end_time - timedelta(days=7)
|
||||
if dataset == "realtime":
|
||||
rows = self._request_rows(cfg)
|
||||
elif dataset == "minute":
|
||||
elif dataset in {"minute", "full_minute"}:
|
||||
override: dict[str, Any] = {}
|
||||
if cfg.asset_type_param:
|
||||
override[cfg.asset_type_param] = "stock"
|
||||
|
||||
@@ -1064,6 +1064,100 @@ def fetch_intraday_universe_increment(
|
||||
return (_normalize_minute(seg), 1)
|
||||
|
||||
|
||||
def _resolve_full_minute_provider(
|
||||
provider_name: str,
|
||||
) -> tuple[object | None, bool, str | None]:
|
||||
"""解析全量分钟生效的自定义源。返回 (provider, should_use_tickflow, error_msg):
|
||||
|
||||
- provider_name == "tickflow" / 未配 full_minute dataset → (None, True, None)
|
||||
- resolver 异常 (registry 损坏 / 插件失效 / 源不存在) → (None, True, str(e))
|
||||
- 成功 → (provider, False, None)
|
||||
|
||||
与 _resolve_minute_provider 同构, 仅数据集名不同。
|
||||
"""
|
||||
if provider_name == "tickflow":
|
||||
return (None, True, None)
|
||||
from app.data_providers import custom as custom_sources
|
||||
try:
|
||||
if not custom_sources.provider_has_dataset(provider_name, "full_minute"):
|
||||
return (None, True, None)
|
||||
provider = custom_sources.get_provider(provider_name)
|
||||
return (provider, False, None)
|
||||
except Exception as e: # noqa: BLE001
|
||||
return (None, True, str(e))
|
||||
|
||||
|
||||
def fetch_intraday_custom_batch(
|
||||
provider: object,
|
||||
provider_name: str,
|
||||
symbols: list[str],
|
||||
) -> tuple[pl.DataFrame, int]:
|
||||
"""自定义源全量分钟修复轮: 当日窗口全市场批量拉取 (不落盘)。
|
||||
|
||||
provider 契约 (见 docs/plugin-development.md):
|
||||
- get_intraday_batch(symbols, count, asset_type) 优先 — 源自管批量端点;
|
||||
- 未实现则回退 get_minute(symbols, 当日窗口) — 复用逐标的分钟K机制,
|
||||
请求数按 chunk 回调统计。
|
||||
帧统一过北京墙钟守卫 (与 _try_custom_minute 同纪律)。
|
||||
返回 (当日分钟K, 请求数); 失败返回空 df 由调用方按空轮处理。
|
||||
"""
|
||||
try:
|
||||
method = getattr(provider, "get_intraday_batch", None)
|
||||
if callable(method):
|
||||
df = method(symbols)
|
||||
requests = 1
|
||||
else:
|
||||
counted = {"requests": 0}
|
||||
|
||||
def _count_requests(cur: int, total: int) -> None:
|
||||
counted["requests"] = max(counted["requests"], int(total))
|
||||
|
||||
end = cn_now()
|
||||
start = end.replace(hour=0, minute=0, second=0, microsecond=0)
|
||||
df = provider.get_minute(
|
||||
symbols, start_time=start, end_time=end,
|
||||
asset_type="stock", freq="1m", on_chunk_done=_count_requests,
|
||||
)
|
||||
requests = counted["requests"]
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("custom full_minute batch via %s failed: %s", provider_name, e)
|
||||
return (pl.DataFrame(), 0)
|
||||
try:
|
||||
df = _enforce_minute_beijing_wallclock(df, source=provider_name)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("custom full_minute datetime 契约校验失败 (%s): %s", provider_name, e)
|
||||
return (pl.DataFrame(), 0)
|
||||
return (df, max(requests, 1))
|
||||
|
||||
|
||||
def fetch_intraday_custom_latest(
|
||||
provider: object,
|
||||
provider_name: str,
|
||||
*,
|
||||
count: int = 3,
|
||||
) -> tuple[pl.DataFrame, int] | None:
|
||||
"""自定义源全量分钟稳态增量轮 (不落盘)。
|
||||
|
||||
provider 可选实现 get_intraday_latest(symbols=None, count) → 每只标的最新
|
||||
count 根分钟K, 尽量单请求/低请求量 (TickFlow 的 intraday.universe 同义)。
|
||||
未实现返回 None — 调用方降级为仅修复轮模式。失败返回 (空 df, 0) 按空轮处理。
|
||||
"""
|
||||
method = getattr(provider, "get_intraday_latest", None)
|
||||
if not callable(method):
|
||||
return None
|
||||
try:
|
||||
df = method(count=count)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("custom full_minute latest via %s failed: %s", provider_name, e)
|
||||
return (pl.DataFrame(), 0)
|
||||
try:
|
||||
df = _enforce_minute_beijing_wallclock(df, source=provider_name)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("custom full_minute datetime 契约校验失败 (%s): %s", provider_name, e)
|
||||
return (pl.DataFrame(), 0)
|
||||
return (df, 1)
|
||||
|
||||
|
||||
def fetch_minute_single(
|
||||
symbol: str,
|
||||
trade_date: date,
|
||||
|
||||
@@ -1,34 +1,41 @@
|
||||
"""盘中分钟K增量落盘服务 (Expert 专有)。
|
||||
"""盘中分钟K增量落盘服务 (全量分钟能力, 按能力路由)。
|
||||
|
||||
两段式拉取 (见 feat/minute-strategy 方案):
|
||||
- 全天修复轮: intraday.batch (日内分时批量) 并发脉冲一次拉全市场当日全部
|
||||
分钟K — 冷启动 (如 10 点才开服务, 补 9:30 起缺口) / 覆盖滞后超阈值 /
|
||||
连续空轮自愈时触发。
|
||||
- 稳态增量轮: intraday.universe 传 CN_Equity_A 标的池, 单请求返回全市场
|
||||
每只最新 3 根 (服务端上限), 靠 _write_minute_partition 的
|
||||
unique(symbol,datetime) 幂等合并滚出全天。
|
||||
- 全天修复轮: 全市场日内分时批量并发脉冲一次拉当日全部分钟K — 冷启动
|
||||
(如 10 点才开服务, 补 9:30 起缺口) / 覆盖滞后超阈值 / 连续空轮自愈时触发。
|
||||
- 稳态增量轮: 传标的池/universe 单请求返回全市场每只最新 3 根, 靠
|
||||
_write_minute_partition 的 unique(symbol,datetime) 幂等合并滚出全天。
|
||||
|
||||
单轮合并写入当日 kline_minute 分区, 供分钟策略 (minute_filter) 读到新鲜数据。
|
||||
|
||||
数据源路由 (像其他能力一样可接入):
|
||||
- 生效源由偏好 full_minute_data_provider 决定 (设置 → 数据源 → 全量分钟):
|
||||
TickFlow (Expert 档, intraday.batch + intraday.universe) 或声明 full_minute
|
||||
数据集的插件/自定义源。
|
||||
- 自定义源契约 (docs/plugin-development.md): get_intraday_batch(symbols) 修复轮
|
||||
(未实现自动回退 get_minute 当日窗口); get_intraday_latest(count) 稳态增量
|
||||
(可选, 未实现自动降级为仅修复轮并放慢节奏至 >=60s)。
|
||||
- 能力门控 Cap.INTRADAY_UNIVERSE: TickFlow 按档位探测, 自定义源声明数据集且
|
||||
被路由时由 policy._augment_custom_sources 补授 — 门控口径对两类源统一。
|
||||
|
||||
设计约束:
|
||||
- Expert 专有: 能力门控 Cap.INTRADAY_UNIVERSE (全量分钟) — 仅 TickFlow
|
||||
Expert 档具备, 天然排他 (自定义分钟源无此能力, 且服务本就让位插件)。
|
||||
- 修复轮并发脉冲: 全市场按 batch_size 分块 (5546/200 = 28 块) 一次打出。
|
||||
任何 60s 滑动窗口至多一个脉冲 (28 < 48 安全 rpm); 单块失败不拖垮整轮,
|
||||
失败块单独重试一次 (见 fetch_intraday_full_market_burst)。
|
||||
- 修复轮并发脉冲: TickFlow 路径全市场按 batch_size 分块 (5546/200 = 28 块)
|
||||
一次打出。任何 60s 滑动窗口至多一个脉冲 (28 < 48 安全 rpm); 单块失败不拖垮
|
||||
整轮, 失败块单独重试一次 (见 fetch_intraday_full_market_burst)。自定义源
|
||||
由 provider 自管批量与限速 (rpm 配置/内部并发), 服务不代限。
|
||||
- 稳态轮单请求: 无脉冲并发, 间隔可低至 3s; 实际节奏 = max(间隔, 单轮完成),
|
||||
服务端响应 ~5s 时自动退化为响应节奏, 不会重叠请求。
|
||||
- 固定节奏: 默认 6s 一轮 (clamp [3, 300]), 不补跑 (missed 轮次直接跳过)。
|
||||
- 固定节奏: 默认 6s 一轮 (clamp [3, 300]), 不补跑 (missed 轮次直接跳过);
|
||||
仅修复轮的自定义源下限抬到 60s (全天批量打不住 6s 节奏)。
|
||||
- 仅连续竞价时段运行 (9:30-11:30 / 13:00-15:00), 午休/收盘自动暂停与恢复;
|
||||
午休后恢复因覆盖滞后会多跑一次修复轮, 幂等无害。
|
||||
- 不与其他分钟能力冲突: 与 盘后分钟同步 (kline.minute.batch) / 分时监控路径
|
||||
分属不同限流池; 落盘走 _write_minute_partition 的 unique(symbol,datetime)
|
||||
合并, 与盘后同步写同一分区安全幂等。
|
||||
- 数据源插件化让位: 配置了自定义分钟源 (minute_data_provider != tickflow) 时
|
||||
服务不启动 — 盘中增量交由插件自管, 本服务不抢占。
|
||||
|
||||
分层: 本模块只做调度/落盘/状态; TickFlow SDK 调用全部在 kline_sync 边界层
|
||||
(fetch_intraday_full_market_burst / fetch_intraday_universe_increment),
|
||||
分层: 本模块只做调度/落盘/状态; 取数全部经 kline_sync 边界层
|
||||
(TickFlow: fetch_intraday_full_market_burst / fetch_intraday_universe_increment;
|
||||
自定义: fetch_intraday_custom_batch / fetch_intraday_custom_latest),
|
||||
保持插件化边界不泄漏。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
@@ -58,6 +65,9 @@ _LOOP_STEP_S = 2.0
|
||||
_REPAIR_LAG_MINUTES = 3.0
|
||||
# 连续空轮达到该次数 → 强制全天修复轮 (自愈 universe 端点持续异常)。
|
||||
_EMPTY_ROUNDS_TO_REPAIR = 2
|
||||
# 仅修复轮的自定义源 (无 get_intraday_latest 廉价增量): 全天批量节奏下限,
|
||||
# 6s 一轮全市场拉取会把绝大多数源打爆。
|
||||
_REPAIR_ONLY_MIN_INTERVAL_S = 60
|
||||
|
||||
|
||||
def _in_continuous_session(now=None) -> bool:
|
||||
@@ -124,8 +134,10 @@ class MinuteRefreshService:
|
||||
def capability_ok(self) -> bool:
|
||||
"""Cap.INTRADAY_UNIVERSE (全量分钟) 存在。能力探测结果缓存在 app.state。
|
||||
|
||||
门控挂在稳态增量的主能力上; 全天修复轮用的 intraday.batch 与其
|
||||
同属 Expert 档 (tiers.yaml), 目前两者必然同时持有。
|
||||
门控挂在稳态增量的主能力上; TickFlow 路径全天修复轮用的 intraday.batch
|
||||
与其同属 Expert 档 (tiers.yaml), 两者必然同时持有。自定义源声明
|
||||
full_minute 数据集且被路由时, 由 policy._augment_custom_sources 补授
|
||||
同一能力键 — 门控口径对两类源统一。
|
||||
"""
|
||||
capset = getattr(self._app_state, "capabilities", None) if self._app_state else None
|
||||
if capset is None:
|
||||
@@ -137,19 +149,54 @@ class MinuteRefreshService:
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
def custom_provider_active(self) -> bool:
|
||||
"""配置了自定义分钟源 → 让位插件, 本服务不启动。"""
|
||||
def active_provider(self) -> str:
|
||||
"""全量分钟生效源名 ("tickflow" 或自定义源名, 偏好 full_minute_data_provider)。"""
|
||||
try:
|
||||
return preferences.get_minute_data_provider() != "tickflow"
|
||||
except Exception:
|
||||
return False
|
||||
return preferences.get_full_minute_data_provider()
|
||||
except Exception: # noqa: BLE001 — 偏好文件异常按 TickFlow 处理
|
||||
return "tickflow"
|
||||
|
||||
def _resolve_custom(self) -> tuple[object | None, str]:
|
||||
"""解析自定义源。返回 (provider_or_None, effective_name):
|
||||
|
||||
- 偏好 tickflow / 源未声明 full_minute 数据集 / 解析异常 → (None, "tickflow")
|
||||
(与 minute 数据集同纪律: 静默降级 TickFlow, 能力门控决定能否真正运行)
|
||||
- 成功 → (provider, name)
|
||||
"""
|
||||
from app.services import kline_sync
|
||||
|
||||
name = self.active_provider()
|
||||
if name == "tickflow":
|
||||
return (None, "tickflow")
|
||||
provider, use_tickflow, err = kline_sync._resolve_full_minute_provider(name)
|
||||
if use_tickflow:
|
||||
if err is not None:
|
||||
logger.warning(
|
||||
"full_minute provider %s 解析失败, 本轮降级 TickFlow: %s", name, err,
|
||||
)
|
||||
return (None, "tickflow")
|
||||
return (provider, name)
|
||||
|
||||
def _custom_supports_increment(self, provider: object) -> bool:
|
||||
"""自定义源是否实现 get_intraday_latest (稳态增量); 否则仅修复轮。"""
|
||||
return callable(getattr(provider, "get_intraday_latest", None))
|
||||
|
||||
def repair_only(self) -> bool:
|
||||
"""当前生效源只能全天修复轮 (无廉价增量端点) — 节奏下限抬到 60s。"""
|
||||
provider, name = self._resolve_custom()
|
||||
return provider is not None and not self._custom_supports_increment(provider)
|
||||
|
||||
def _effective_interval(self) -> int:
|
||||
"""本轮间隔: 偏好值; 仅修复轮的自定义源下限 60s (全天批量打不住 6s 节奏)。"""
|
||||
interval = preferences.get_minute_refresh_interval()
|
||||
if self.repair_only():
|
||||
return max(interval, _REPAIR_ONLY_MIN_INTERVAL_S)
|
||||
return interval
|
||||
|
||||
def _gate_reason(self) -> str | None:
|
||||
"""返回本轮不执行的原因 (None = 放行)。"""
|
||||
if not preferences.get_minute_refresh_enabled():
|
||||
return "disabled"
|
||||
if self.custom_provider_active():
|
||||
return "custom_minute_provider"
|
||||
if not self.capability_ok():
|
||||
return "capability"
|
||||
if not _in_continuous_session():
|
||||
@@ -171,7 +218,7 @@ class MinuteRefreshService:
|
||||
try:
|
||||
reason = self._gate_reason()
|
||||
if reason is None:
|
||||
interval = preferences.get_minute_refresh_interval()
|
||||
interval = self._effective_interval()
|
||||
started = time.time()
|
||||
self._run_round()
|
||||
# 固定节奏: 下一轮 = max(本轮起点+间隔, 本轮完成), 不补跑
|
||||
@@ -234,13 +281,21 @@ class MinuteRefreshService:
|
||||
from app.services import kline_sync
|
||||
|
||||
t0 = time.perf_counter()
|
||||
custom, provider_name = self._resolve_custom()
|
||||
mode = self._select_mode()
|
||||
if mode == "increment" and custom is not None and not self._custom_supports_increment(custom):
|
||||
# 无廉价增量端点的源: 增量轮退化为修复轮 (全天批量幂等覆盖, 数据不丢)
|
||||
mode = "full"
|
||||
mode_label = "增量" if mode == "increment" else "全天修复"
|
||||
with self._round_lock:
|
||||
if mode == "increment":
|
||||
# fetch 计时只覆盖网络取数; full 分支的 universe 维表读取不计入
|
||||
fetch_started = time.perf_counter()
|
||||
df, requests = kline_sync.fetch_intraday_universe_increment()
|
||||
if custom is not None:
|
||||
latest = kline_sync.fetch_intraday_custom_latest(custom, provider_name)
|
||||
df, requests = latest if latest is not None else (pl.DataFrame(), 0)
|
||||
else:
|
||||
df, requests = kline_sync.fetch_intraday_universe_increment()
|
||||
self._state.last_symbols = (
|
||||
df["symbol"].n_unique() if not df.is_empty() else 0
|
||||
)
|
||||
@@ -251,9 +306,14 @@ class MinuteRefreshService:
|
||||
self._state.last_error = "empty universe (instruments 未加载)"
|
||||
logger.warning("全量分钟[%s] 本轮中止: 标的池为空 (instruments 未加载)", mode_label)
|
||||
return
|
||||
capset = getattr(self._app_state, "capabilities", None) if self._app_state else None
|
||||
fetch_started = time.perf_counter()
|
||||
df, requests = kline_sync.fetch_intraday_full_market_burst(symbols, capset)
|
||||
if custom is not None:
|
||||
df, requests = kline_sync.fetch_intraday_custom_batch(
|
||||
custom, provider_name, symbols,
|
||||
)
|
||||
else:
|
||||
capset = getattr(self._app_state, "capabilities", None) if self._app_state else None
|
||||
df, requests = kline_sync.fetch_intraday_full_market_burst(symbols, capset)
|
||||
fetch_ms = (time.perf_counter() - fetch_started) * 1000
|
||||
self._state.last_requests = requests
|
||||
if df.is_empty():
|
||||
@@ -326,9 +386,11 @@ class MinuteRefreshService:
|
||||
"enabled": enabled,
|
||||
"running": running,
|
||||
"healthy": self.is_healthy(),
|
||||
"provider": self.active_provider(),
|
||||
"provider_effective": self._resolve_custom()[1],
|
||||
"repair_only": self.repair_only(),
|
||||
"interval_seconds": preferences.get_minute_refresh_interval(),
|
||||
"capability_ok": self.capability_ok(),
|
||||
"custom_provider_active": self.custom_provider_active(),
|
||||
"in_trading_hours": _in_continuous_session(),
|
||||
"gate_reason": gate if (enabled and running) else (gate or "disabled"),
|
||||
"rounds": self._state.rounds,
|
||||
@@ -343,8 +405,8 @@ class MinuteRefreshService:
|
||||
}
|
||||
|
||||
def trigger_manual_round(self) -> dict[str, Any]:
|
||||
"""手动触发一轮 (无视时段门控, 但仍受能力/插件门控); 供状态页「立即刷新」。"""
|
||||
if self.custom_provider_active() or not self.capability_ok():
|
||||
"""手动触发一轮 (无视时段门控, 但仍受能力门控); 供状态页「立即刷新」。"""
|
||||
if not self.capability_ok():
|
||||
return {"ok": False, "reason": self._gate_reason() or "capability"}
|
||||
threading.Thread(target=self._run_round, daemon=True, name="minute-refresh-manual").start()
|
||||
return {"ok": True}
|
||||
|
||||
@@ -285,6 +285,11 @@ def get_minute_data_provider() -> str:
|
||||
return provider if provider in _allowed_data_providers() else "tickflow"
|
||||
|
||||
|
||||
def get_full_minute_data_provider() -> str:
|
||||
provider = str(load().get("full_minute_data_provider", "tickflow") or "tickflow").lower()
|
||||
return provider if provider in _allowed_data_providers() else "tickflow"
|
||||
|
||||
|
||||
def get_depth5_data_provider() -> str:
|
||||
provider = str(load().get("depth5_data_provider", "tickflow") or "tickflow").lower()
|
||||
return provider if provider in _allowed_data_providers() else "tickflow"
|
||||
|
||||
@@ -311,6 +311,7 @@ _DATASET_CAP_MAP: tuple[tuple[str, Cap], ...] = (
|
||||
("adj_factor", Cap.ADJ_FACTOR),
|
||||
("minute", Cap.KLINE_MINUTE_BATCH),
|
||||
("financial", Cap.FINANCIAL),
|
||||
("full_minute", Cap.INTRADAY_UNIVERSE),
|
||||
)
|
||||
|
||||
|
||||
@@ -327,6 +328,7 @@ def _augment_custom_sources(capset: CapabilitySet) -> None:
|
||||
"adj_factor": adj_provider,
|
||||
"minute": preferences.get_minute_data_provider(),
|
||||
"financial": preferences.get_financial_provider(),
|
||||
"full_minute": preferences.get_full_minute_data_provider(),
|
||||
}
|
||||
for dataset, cap in _DATASET_CAP_MAP:
|
||||
provider = active_providers[dataset]
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
"""能力标准统一: 自定义/插件数据源能力增广回归测试。
|
||||
|
||||
对应 _augment_custom_sources 的数据集→能力映射 (daily/adj_factor/minute/financial):
|
||||
对应 _augment_custom_sources 的数据集→能力映射 (daily/adj_factor/minute/financial/full_minute):
|
||||
某数据集的当前 provider 非 tickflow 且声明了该数据集 → grant 对应能力;
|
||||
取数路由仍按 preferences 分流, 不会误调 TickFlow。
|
||||
"""
|
||||
@@ -13,13 +13,14 @@ from app.tickflow.policy import _augment_custom_sources
|
||||
|
||||
|
||||
def _set_providers(monkeypatch, *, daily="tickflow", adj="tickflow",
|
||||
minute="tickflow", financial="tickflow") -> None:
|
||||
minute="tickflow", financial="tickflow", full_minute="tickflow") -> None:
|
||||
"""mock preferences 各数据集 provider getter。"""
|
||||
from app.services import preferences
|
||||
monkeypatch.setattr(preferences, "get_daily_data_provider", lambda: daily)
|
||||
monkeypatch.setattr(preferences, "get_adj_factor_provider", lambda: adj)
|
||||
monkeypatch.setattr(preferences, "get_minute_data_provider", lambda: minute)
|
||||
monkeypatch.setattr(preferences, "get_financial_provider", lambda: financial)
|
||||
monkeypatch.setattr(preferences, "get_full_minute_data_provider", lambda: full_minute)
|
||||
|
||||
|
||||
def _set_datasets(monkeypatch, datasets: set[str]) -> None:
|
||||
@@ -51,6 +52,27 @@ def test_adj_custom_source_grants_adj_factor(monkeypatch):
|
||||
assert capset.has(Cap.ADJ_FACTOR)
|
||||
|
||||
|
||||
def test_full_minute_custom_source_grants_intraday_universe(monkeypatch):
|
||||
"""全量分钟: 声明 full_minute 数据集且被路由 → 补授 INTRADAY_UNIVERSE,
|
||||
minute_refresh 服务门控与 TickFlow Expert 口径统一。"""
|
||||
_set_providers(monkeypatch, full_minute="mock_src")
|
||||
_set_datasets(monkeypatch, {"full_minute"})
|
||||
capset = CapabilitySet()
|
||||
_augment_custom_sources(capset)
|
||||
assert capset.has(Cap.INTRADAY_UNIVERSE)
|
||||
# 声明了 full_minute 不等于声明 minute → 不补逐标的分钟K能力
|
||||
assert not capset.has(Cap.KLINE_MINUTE_BATCH)
|
||||
|
||||
|
||||
def test_full_minute_dataset_without_routing_not_granted(monkeypatch):
|
||||
"""源声明了 full_minute 但路由仍是 tickflow → 不增广 (TickFlow 档位自决)。"""
|
||||
_set_providers(monkeypatch, full_minute="tickflow")
|
||||
_set_datasets(monkeypatch, {"full_minute"})
|
||||
capset = CapabilitySet()
|
||||
_augment_custom_sources(capset)
|
||||
assert not capset.has(Cap.INTRADAY_UNIVERSE)
|
||||
|
||||
|
||||
def test_minute_custom_source_grants_minute_batch(monkeypatch):
|
||||
"""原有 minute 增广行为保持。"""
|
||||
_set_providers(monkeypatch, minute="mock_src")
|
||||
|
||||
@@ -17,6 +17,7 @@ DEFAULT_CURRENT = {
|
||||
"daily_data_provider": "tickflow",
|
||||
"adj_factor_provider": "tickflow",
|
||||
"minute_data_provider": "tickflow",
|
||||
"full_minute_data_provider": "tickflow",
|
||||
"depth5_data_provider": "tickflow",
|
||||
"realtime_data_provider": "tickflow",
|
||||
"financial_data_provider": "tickflow",
|
||||
@@ -42,7 +43,7 @@ def test_registry_covers_all_routing_fields():
|
||||
"realtime", "daily", "minute", "full_minute", "depth5", "adj_factor", "financial",
|
||||
}
|
||||
full_minute = next(c for c in CAPABILITY_REGISTRY if c["id"] == "full_minute")
|
||||
assert full_minute["field"] is None
|
||||
assert full_minute["field"] == "full_minute_data_provider"
|
||||
assert full_minute["tf_tier"] == "expert"
|
||||
for cap in CAPABILITY_REGISTRY:
|
||||
assert cap["default"] == "tickflow"
|
||||
@@ -237,20 +238,33 @@ def test_unknown_current_display_falls_back_to_name(monkeypatch):
|
||||
assert caps["realtime"]["effective_display"] == "ghost"
|
||||
|
||||
|
||||
def test_full_minute_row_is_non_routable_expert_only(monkeypatch):
|
||||
"""全量分钟行: field=None 不可路由, 生效源恒为 TickFlow, 按 expert 档判定可用。"""
|
||||
_fake_sources(monkeypatch, [])
|
||||
# 即使有插件声明别的数据集也不会成为全量分钟候选 (契约不开放该数据集)
|
||||
def test_full_minute_routable_like_other_capabilities(monkeypatch):
|
||||
"""全量分钟行: 与其他能力同样可路由 — 声明 full_minute 数据集的源进候选,
|
||||
路由到它则 usable=True (TickFlow 档位不足也不拦); 未路由且档位不足才不可用。"""
|
||||
_fake_sources(
|
||||
monkeypatch,
|
||||
[],
|
||||
[{"name": "myfm", "display_name": "MyFM", "datasets": ["full_minute"]}],
|
||||
)
|
||||
caps = _by_id(build_capability_matrix(dict(DEFAULT_CURRENT), tickflow_tier="expert"))
|
||||
fm = caps["full_minute"]
|
||||
assert fm["field"] is None
|
||||
assert [c["name"] for c in fm["candidates"]] == ["tickflow"]
|
||||
assert fm["field"] == "full_minute_data_provider"
|
||||
assert [c["name"] for c in fm["candidates"]] == ["tickflow", "myfm"]
|
||||
assert fm["usable"] is True
|
||||
assert fm["tf_available"] is True
|
||||
assert fm["effective"] == "tickflow"
|
||||
|
||||
caps_pro = _by_id(build_capability_matrix(dict(DEFAULT_CURRENT), tickflow_tier="pro"))
|
||||
# TickFlow 档位不足 (pro) 但路由到声明该数据集的自定义源 → 同样可用
|
||||
routed = dict(DEFAULT_CURRENT, full_minute_data_provider="myfm")
|
||||
caps_pro = _by_id(build_capability_matrix(routed, tickflow_tier="pro"))
|
||||
fm_pro = caps_pro["full_minute"]
|
||||
assert fm_pro["candidates"] == []
|
||||
assert fm_pro["usable"] is False
|
||||
assert [c["name"] for c in fm_pro["candidates"]] == ["myfm"]
|
||||
assert fm_pro["usable"] is True
|
||||
assert fm_pro["tf_available"] is False
|
||||
assert fm_pro["effective"] == "myfm"
|
||||
|
||||
# 档位不足且未路由 → 不可用
|
||||
caps_pro_default = _by_id(build_capability_matrix(dict(DEFAULT_CURRENT), tickflow_tier="pro"))
|
||||
fm_default = caps_pro_default["full_minute"]
|
||||
assert [c["name"] for c in fm_default["candidates"]] == ["myfm"]
|
||||
assert fm_default["usable"] is False
|
||||
|
||||
@@ -2,13 +2,15 @@
|
||||
|
||||
覆盖:
|
||||
- 连续竞价时段判定 (含边界)
|
||||
- 门控链: 开关关闭 / 自定义分钟源让位 / 能力缺失 / 时段外 / 放行
|
||||
- 门控链: 开关关闭 / 能力缺失 / 时段外 / 放行 (自定义源与 TickFlow 统一口径)
|
||||
- 数据源路由: full_minute 偏好 → 自定义源 (get_intraday_batch / get_intraday_latest /
|
||||
get_minute 回退) 或降级 TickFlow; 仅修复轮源 60s 节奏下限
|
||||
- 单轮: mock 边界层脉冲 + 落盘, 校验状态字段与 universe 来源
|
||||
- 偏好读写: 默认关闭、间隔 clamp [60, 300]
|
||||
- 偏好读写: 默认关闭、间隔 clamp [3, 120]
|
||||
- API: /minute-refresh/status 无服务时 available=false
|
||||
|
||||
不发起真实网络请求: fetch_intraday_full_market_burst 与 _write_minute_partition
|
||||
均 monkeypatch 替换。
|
||||
不发起真实网络请求: TickFlow 侧边界函数与 _write_minute_partition 均
|
||||
monkeypatch 替换; 自定义源侧用内存 fake provider 走真实边界包装。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
@@ -70,12 +72,20 @@ def test_continuous_session_boundaries():
|
||||
# ── 门控链 ──────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _svc(tmp_path, monkeypatch, *, enabled=True, custom_provider=False, capability=True, in_hours=True):
|
||||
def _svc(tmp_path, monkeypatch, *, enabled=True, capability=True, in_hours=True,
|
||||
full_minute_provider="tickflow", custom=None):
|
||||
"""custom 非 None 时模拟「full_minute 路由到声明该数据集的自定义源」:
|
||||
resolver 直接返回 fake provider (注册表在测试环境未加载)。"""
|
||||
_isolated_prefs(tmp_path, monkeypatch)
|
||||
preferences.save({"minute_refresh_enabled": enabled})
|
||||
if custom_provider:
|
||||
# 模拟已注册的自定义分钟源 (真实注册表在测试环境未加载)
|
||||
monkeypatch.setattr(preferences, "get_minute_data_provider", lambda: "a-stock-data")
|
||||
monkeypatch.setattr(
|
||||
preferences, "get_full_minute_data_provider", lambda: full_minute_provider,
|
||||
)
|
||||
if custom is not None:
|
||||
monkeypatch.setattr(
|
||||
"app.services.kline_sync._resolve_full_minute_provider",
|
||||
lambda name: (custom, False, None),
|
||||
)
|
||||
svc = MinuteRefreshService(_FakeRepo(["600000.SH"]))
|
||||
svc.set_app_state(_FakeAppState(capability))
|
||||
monkeypatch.setattr(minute_refresh, "_in_continuous_session", lambda now=None: in_hours)
|
||||
@@ -90,9 +100,15 @@ def test_gate_disabled(tmp_path, monkeypatch):
|
||||
assert _svc(tmp_path, monkeypatch, enabled=False)._gate_reason() == "disabled"
|
||||
|
||||
|
||||
def test_gate_custom_provider_yields(tmp_path, monkeypatch):
|
||||
svc = _svc(tmp_path, monkeypatch, custom_provider=True)
|
||||
assert svc._gate_reason() == "custom_minute_provider"
|
||||
def test_gate_full_minute_custom_provider_runs(tmp_path, monkeypatch):
|
||||
"""全量分钟路由到自定义源: 不再让位 — 能力口径统一 (capset 增广后放行)。"""
|
||||
class _P: # noqa: D401 — fake provider
|
||||
def get_intraday_batch(self, symbols, count=300, asset_type="stock"):
|
||||
return pl.DataFrame()
|
||||
svc = _svc(tmp_path, monkeypatch, full_minute_provider="myfm", custom=_P())
|
||||
assert svc._gate_reason() is None
|
||||
assert svc.active_provider() == "myfm"
|
||||
assert svc.status()["provider"] == "myfm"
|
||||
|
||||
|
||||
def test_gate_capability_missing(tmp_path, monkeypatch):
|
||||
@@ -108,7 +124,7 @@ def test_gate_outside_trading_hours(tmp_path, monkeypatch):
|
||||
def test_gate_pass(tmp_path, monkeypatch):
|
||||
svc = _svc(tmp_path, monkeypatch)
|
||||
assert svc._gate_reason() is None
|
||||
assert svc.capability_ok() and not svc.custom_provider_active()
|
||||
assert svc.capability_ok() and svc.active_provider() == "tickflow"
|
||||
|
||||
|
||||
# ── 单轮 ────────────────────────────────────────────────────────────
|
||||
@@ -253,6 +269,131 @@ def test_consecutive_empty_rounds_escalate_to_full(tmp_path, monkeypatch):
|
||||
assert svc.status()["last_mode"] == "full"
|
||||
|
||||
|
||||
# ── 自定义源 (full_minute 路由) 轮次 ────────────────────────────────
|
||||
|
||||
|
||||
class _FakeCustomProvider:
|
||||
"""内存 fake: 按 flags 暴露 get_intraday_batch / get_intraday_latest / get_minute。"""
|
||||
|
||||
def __init__(self, *, batch=False, latest=False, minute=False):
|
||||
self.calls: dict = {}
|
||||
# 未启用的方法置 None (实例属性遮蔽类方法) — getattr 返回 None → 边界
|
||||
# 包装按"未实现"处理, 与真实 provider 只暴露部分方法的形状一致
|
||||
if not batch:
|
||||
self.get_intraday_batch = None
|
||||
if not latest:
|
||||
self.get_intraday_latest = None
|
||||
if not minute:
|
||||
self.get_minute = None
|
||||
|
||||
def get_intraday_batch(self, symbols, count=300, asset_type="stock"):
|
||||
self.calls["batch"] = list(symbols)
|
||||
return _full_df()
|
||||
|
||||
def get_intraday_latest(self, symbols=None, count=3):
|
||||
self.calls["latest"] = count
|
||||
return _inc_df()
|
||||
|
||||
def get_minute(self, symbols, start_time=None, end_time=None,
|
||||
asset_type="stock", freq="1m", on_chunk_done=None):
|
||||
self.calls["minute"] = {
|
||||
"symbols": list(symbols),
|
||||
"start": start_time, "end": end_time,
|
||||
}
|
||||
if on_chunk_done:
|
||||
on_chunk_done(1, 4)
|
||||
return _full_df()
|
||||
|
||||
|
||||
def _patch_write(monkeypatch) -> dict:
|
||||
calls = {"rows": None}
|
||||
monkeypatch.setattr(
|
||||
"app.services.kline_sync._write_minute_partition",
|
||||
lambda df, minute_dir: calls.update(rows=df.height) or df.height,
|
||||
)
|
||||
return calls
|
||||
|
||||
|
||||
def test_custom_provider_full_round_via_intraday_batch(tmp_path, monkeypatch):
|
||||
"""修复轮走 provider.get_intraday_batch (真实边界包装含时区守卫)。"""
|
||||
provider = _FakeCustomProvider(batch=True)
|
||||
svc = _svc(tmp_path, monkeypatch, full_minute_provider="myfm", custom=provider)
|
||||
monkeypatch.setattr(
|
||||
minute_refresh.MinuteRefreshService, "_today_coverage_lag_minutes",
|
||||
lambda self: None, # 冷启动 → 全天修复
|
||||
)
|
||||
write_calls = _patch_write(monkeypatch)
|
||||
svc._run_round()
|
||||
assert provider.calls["batch"] == ["600000.SH"]
|
||||
assert write_calls["rows"] == 2
|
||||
st = svc.status()
|
||||
assert st["last_mode"] == "full"
|
||||
assert st["provider"] == "myfm"
|
||||
assert st["provider_effective"] == "myfm"
|
||||
|
||||
|
||||
def test_custom_provider_increment_round_via_latest(tmp_path, monkeypatch):
|
||||
"""实现 get_intraday_latest 的源: 覆盖新鲜时走稳态增量, 不打全天批量。"""
|
||||
provider = _FakeCustomProvider(batch=True, latest=True)
|
||||
svc = _svc(tmp_path, monkeypatch, full_minute_provider="myfm", custom=provider)
|
||||
monkeypatch.setattr(
|
||||
minute_refresh.MinuteRefreshService, "_today_coverage_lag_minutes",
|
||||
lambda self: 0.2,
|
||||
)
|
||||
_patch_write(monkeypatch)
|
||||
svc._run_round()
|
||||
assert provider.calls == {"latest": 3}
|
||||
st = svc.status()
|
||||
assert st["last_mode"] == "increment"
|
||||
assert st["last_requests"] == 1
|
||||
assert st["repair_only"] is False
|
||||
|
||||
|
||||
def test_custom_provider_latest_missing_forces_full_with_minute_fallback(tmp_path, monkeypatch):
|
||||
"""无 get_intraday_latest: 增量轮退化为修复轮, 批量未实现时回退 get_minute
|
||||
(当日窗口), 节奏下限抬到 60s。"""
|
||||
provider = _FakeCustomProvider(minute=True)
|
||||
svc = _svc(tmp_path, monkeypatch, full_minute_provider="myfm", custom=provider)
|
||||
monkeypatch.setattr(
|
||||
minute_refresh.MinuteRefreshService, "_today_coverage_lag_minutes",
|
||||
lambda self: 0.2, # 覆盖新鲜 → 本应增量
|
||||
)
|
||||
_patch_write(monkeypatch)
|
||||
svc._run_round()
|
||||
# 增量退化为修复: 走 get_minute 回退 (chunk 回调统计请求数)
|
||||
assert "minute" in provider.calls
|
||||
assert provider.calls["minute"]["symbols"] == ["600000.SH"]
|
||||
assert svc.status()["last_mode"] == "full"
|
||||
assert svc.status()["last_requests"] == 4
|
||||
assert svc.repair_only() is True
|
||||
assert svc._effective_interval() == 60
|
||||
|
||||
|
||||
def test_custom_provider_unresolved_degrades_to_tickflow(tmp_path, monkeypatch):
|
||||
"""路由指向未声明 full_minute 数据集的源 (真实 resolver 判定) → 降级
|
||||
TickFlow 路径, 门控口径不变。"""
|
||||
svc = _svc(tmp_path, monkeypatch, full_minute_provider="not-registered")
|
||||
monkeypatch.setattr(
|
||||
minute_refresh.MinuteRefreshService, "_today_coverage_lag_minutes",
|
||||
lambda self: 0.2,
|
||||
)
|
||||
_patch_write(monkeypatch)
|
||||
modes: list[str] = []
|
||||
monkeypatch.setattr(
|
||||
"app.services.kline_sync.fetch_intraday_universe_increment",
|
||||
lambda *a, **k: (modes.append("tf-inc"), (_inc_df(), 1))[1],
|
||||
)
|
||||
monkeypatch.setattr(
|
||||
"app.services.kline_sync.fetch_intraday_full_market_burst",
|
||||
lambda symbols, capset, *, count=300: (modes.append("tf-full"), (_full_df(), 28))[1],
|
||||
)
|
||||
svc._run_round()
|
||||
assert modes == ["tf-inc"]
|
||||
st = svc.status()
|
||||
assert st["provider_effective"] == "tickflow"
|
||||
assert st["last_mode"] == "increment"
|
||||
|
||||
|
||||
def test_status_reports_gate_reason_when_stopped(tmp_path, monkeypatch):
|
||||
svc = _svc(tmp_path, monkeypatch, enabled=False)
|
||||
st = svc.status()
|
||||
|
||||
@@ -205,7 +205,6 @@ def _minute_service(monkeypatch):
|
||||
"app.services.minute_refresh.preferences.get_minute_refresh_enabled",
|
||||
lambda: True,
|
||||
)
|
||||
monkeypatch.setattr(svc, "custom_provider_active", lambda: False)
|
||||
monkeypatch.setattr(svc, "capability_ok", lambda: True)
|
||||
return svc
|
||||
|
||||
|
||||
@@ -33,16 +33,19 @@ TickFlow 是内置默认数据源;同时支持插件化接入第三方数据源(
|
||||
|
||||
### 全量分钟 (full_minute)
|
||||
|
||||
「全量分钟」是一项**独立能力**(能力键 `full_minute`,探测名 `intraday.universe`),不是某个档位的封闭功能:盘中把全市场当日 1 分钟 K 持续增量落盘到本地 `data/kline_minute/` 当日分区,分钟策略(`minute_filter`)与分时视图即可读到新鲜数据。**当前接入方式为 TickFlow Expert** — 配置 Expert 档 Key 即可使用。
|
||||
「全量分钟」是一项**独立能力**(能力键 `full_minute`,探测名 `intraday.universe`),与其他能力同样**可路由**:盘中把全市场当日 1 分钟 K 持续增量落盘到本地 `data/kline_minute/` 当日分区,分钟策略(`minute_filter`)与分时视图即可读到新鲜数据。接入方式二选一:
|
||||
|
||||
- **TickFlow Expert**:配置 Expert 档 Key,零配置即用(修复轮 `intraday.batch` + 稳态 `intraday.universe` 单请求增量)
|
||||
- **插件/自定义源**:声明 `full_minute` 数据集并在 **设置 → 数据源 → 全量分钟** 路由到该源 — Python 插件实现 `get_intraday_batch`(必需)/`get_intraday_latest`(可选,未实现自动降级仅修复轮、节奏下限 60s);YAML 声明式源数据集配置与 `minute` 同形(仅修复轮语义)。契约细节见 [plugin-development.md](./plugin-development.md) 与 [custom-data-source.md](./custom-data-source.md)
|
||||
|
||||
接入步骤:
|
||||
|
||||
1. **设置 → 凭据与能力** 配置 TickFlow API Key(Expert 档),点「重新检测」,能力列表出现「全量分钟」
|
||||
1. **设置 → 凭据与能力**(TickFlow 路径)配置 API Key(Expert 档),或在 **设置 → 数据源** 声明/安装提供 `full_minute` 的源并路由;点「重新检测」后能力列表出现「全量分钟」
|
||||
2. **开启实时行情**后落盘服务自动启动;仅连续竞价时段(9:30–11:30 / 13:00–15:00)运行,午休/收盘自动暂停与恢复
|
||||
3. 冷启动(如 10 点才开服务)自动触发**全天修复轮**,一次并发脉冲补齐 9:30 起的全部缺口;稳态走**增量轮**(默认 6 秒一轮,可配 3–120 秒),全市场每只取最新 3 根幂等合并滚出全天
|
||||
3. 冷启动(如 10 点才开服务)自动触发**全天修复轮**,一次批量补齐 9:30 起的全部缺口;稳态走**增量轮**(默认 6 秒一轮,可配 3–120 秒),幂等合并滚出全天
|
||||
4. 与盘后分钟同步写同一分区(`unique(symbol, datetime)` 幂等合并),互不冲突
|
||||
|
||||
说明:标的池为 A 股股票(CN_Equity_A),ETF 不在内(分时走批量补拉路径);覆盖滞后超阈值或连续空轮会自动再跑修复轮自愈。该数据集暂未纳入插件 loader 白名单;但配置了自定义分钟源(`minute_data_provider ≠ tickflow`)时,内置服务自动让位不启动,盘中分时由你的数据源按需供给(见 [plugin-development.md](./plugin-development.md))。
|
||||
说明:标的池为 A 股股票(CN_Equity_A),ETF 不在内(分时走批量补拉路径);覆盖滞后超阈值或连续空轮会自动再跑修复轮自愈。
|
||||
|
||||
---
|
||||
|
||||
|
||||
@@ -4,7 +4,7 @@
|
||||
|
||||
## 支持范围
|
||||
|
||||
当前自定义源支持五类数据:
|
||||
当前自定义源支持六类数据:
|
||||
|
||||
| 数据集 | 配置名 | 说明 |
|
||||
| --- | --- | --- |
|
||||
@@ -12,10 +12,15 @@
|
||||
| 除权因子 | `adj_factor` | 批量返回一组股票的复权因子 |
|
||||
| 实时行情 | `realtime` | 返回全市场快照,用于盘中 enriched 增量计算 |
|
||||
| 分钟K | `minute` | 返回 1m 分钟K(需映射出 symbol / datetime / OHLC / 量额) |
|
||||
| 全量分钟 | `full_minute` | 与 `minute` 同形;声明后可被路由为「全量分钟」生效源,内置服务盘中按当日窗口全市场批量落盘(仅修复轮语义,节奏下限 60s) |
|
||||
| 财务数据 | `financial` | 一个配置覆盖全部财务表,请求时把表名作为参数传给上游;字段由数据源决定,仅需映射出 symbol |
|
||||
|
||||
深度盘口(depth5)暂无数据集契约,仍由 TickFlow 提供。
|
||||
|
||||
`full_minute` 声明式源只提供修复轮(当日窗口批量);廉价增量端点
|
||||
(`get_intraday_latest`)是 Python 插件契约,见
|
||||
[plugin-development.md](./plugin-development.md)。
|
||||
|
||||
## 配置位置
|
||||
|
||||
把 YAML 放到运行数据目录下:
|
||||
@@ -294,7 +299,7 @@ cp docs/examples/custom-data-source/mock_source.yaml data/data_sources/mock_sour
|
||||
# 上游若返回百分数值 (3.66 表示 3.66%), 在 realtime 数据集声明 pct_unit: percent,
|
||||
# 不要依赖数值自动识别; 逐列转换也可用 transforms: turnover_rate: "value / 100"
|
||||
|
||||
分钟K (minute):
|
||||
分钟K (minute) 与 全量分钟 (full_minute, 字段同 minute):
|
||||
symbol = 股票代码
|
||||
# datetime 必须是北京时间墙钟 (如 2026-08-28 09:35:00), 不要返回 UTC;
|
||||
# 入口守卫会自动纠偏 UTC 特征帧, 但契约仍要求源头写对
|
||||
|
||||
+1
-1
@@ -134,7 +134,7 @@
|
||||
|
||||
### 全量分钟(盘中增量落盘)
|
||||
|
||||
「全量分钟」是独立能力(`full_minute`):盘中将全市场当日 1 分钟 K 持续增量落盘 `data/kline_minute/`,分钟策略与分时视图即取即用。**当前通过 TickFlow Expert 接入** — 配 Key 检测出能力后,开启实时行情即自动运行(冷启动全天修复轮 + 稳态增量轮);配置自定义分钟源时内置服务自动让位。
|
||||
「全量分钟」是独立能力(`full_minute`),**与其他能力同样可路由**:盘中将全市场当日 1 分钟 K 持续增量落盘 `data/kline_minute/`,分钟策略与分时视图即取即用。接入方式:TickFlow Expert,或声明 `full_minute` 数据集的插件/自定义源(设置 → 数据源 → 全量分钟 路由);配好能力后开启实时行情即自动运行(冷启动全天修复轮 + 稳态增量轮)。
|
||||
|
||||
详见 [configuration.md → 全量分钟](./configuration.md#全量分钟-full_minute)。
|
||||
|
||||
|
||||
@@ -149,6 +149,14 @@ class MyProvider:
|
||||
on_chunk_done=None, freq="1m") -> pl.DataFrame:
|
||||
"""分钟K: [symbol, datetime(北京墙钟), open, high, low, close, volume, amount]"""
|
||||
|
||||
def get_intraday_batch(self, symbols, count=300, asset_type="stock") -> pl.DataFrame:
|
||||
"""(声明 full_minute 数据集时实现) 全量分钟修复轮: 给定标的当日 1 分钟K,
|
||||
canonical 8 列同 get_minute; 内部自行分块/限速。"""
|
||||
|
||||
def get_intraday_latest(self, symbols=None, count=3) -> pl.DataFrame:
|
||||
"""(可选, full_minute 稳态增量轮) 尽量单请求返回全市场每只最新 count 根;
|
||||
未实现则服务降级为仅修复轮 (节奏下限 60s)。"""
|
||||
|
||||
def get_realtime(self) -> list[dict]:
|
||||
"""全市场实时快照 → list[dict]。失败软返回 [], 不抛异常(不阻断轮询线程)。"""
|
||||
|
||||
@@ -184,10 +192,25 @@ class MyProvider:
|
||||
深历史(TickFlow 基准)。浅源(如 stock-sdk 免费分时仅保留最近 5 个交易日)声明后,
|
||||
个股分时档位自动收窄为可行选项并默认 5 日,深源默认 20 日。
|
||||
|
||||
> **全量分钟与插件源的关系**:「全量分钟」(`intraday.universe`,盘中全市场分钟增量
|
||||
> 落盘)是独立能力,当前经 TickFlow Expert 接入、由内置服务(`minute_refresh`)执行,
|
||||
> 该数据集暂未纳入插件 loader 白名单。但配置了自定义分钟源(`minute_data_provider
|
||||
> != tickflow`)时,内置服务自动让位不启动,盘中分时由你的源经 `get_minute` 按需供给。
|
||||
> **全量分钟 (full_minute) 数据集契约**:声明 `full_minute` 数据集并把
|
||||
> `full_minute_data_provider` 路由到你的源,即接入「全量分钟」能力(盘中全市场
|
||||
> 当日分钟K增量落盘,由内置服务 `minute_refresh` 调度,与 TickFlow Expert 同一
|
||||
> 能力键 `intraday.universe`)。需实现:
|
||||
>
|
||||
> - `get_intraday_batch(symbols, count=300, asset_type="stock") -> pl.DataFrame`
|
||||
> — **必须**(或已有 `get_minute` 自动回退,但强烈建议实现批量端点)。
|
||||
> 返回给定标的当日 1 分钟K,canonical 8 列
|
||||
> `[symbol, datetime(北京墙钟 naive), open, high, low, close, volume, amount]`,
|
||||
> 内部自行分块/限速。服务在冷启动、覆盖断档、连续空轮时调用(修复轮)。
|
||||
> - `get_intraday_latest(symbols=None, count=3) -> pl.DataFrame` — **可选**,
|
||||
> 稳态增量轮专用:尽量单请求返回全市场每只标的最新 `count` 根。未实现时
|
||||
> 服务自动降级为仅修复轮,节奏下限抬到 60s(全天批量打不住 6s 节奏)。
|
||||
>
|
||||
> 两个方法的返回帧都过 `_enforce_minute_beijing_wallclock` 时区守卫(与
|
||||
> `get_minute` 同纪律);失败抛异常或返回空 df 均按空轮处理,连续空轮触发
|
||||
> 修复轮自愈。声明方式:插件在 `plugin.yaml` 的 `datasets:` 列表加入
|
||||
> `full_minute`。YAML 声明式源同样支持(数据集配置与 `minute` 同形,仅提供
|
||||
> 修复轮语义,见 [custom-data-source.md](./custom-data-source.md))。
|
||||
|
||||
### 异常语义
|
||||
|
||||
|
||||
@@ -1485,6 +1485,7 @@ export type ProviderField =
|
||||
| 'daily_data_provider'
|
||||
| 'adj_factor_provider'
|
||||
| 'minute_data_provider'
|
||||
| 'full_minute_data_provider'
|
||||
| 'depth5_data_provider'
|
||||
| 'realtime_data_provider'
|
||||
| 'financial_data_provider'
|
||||
@@ -1604,6 +1605,8 @@ export interface Preferences {
|
||||
daily_data_provider?: string
|
||||
adj_factor_provider?: string
|
||||
minute_data_provider?: string
|
||||
/** 全量分钟 (盘中全市场分钟落盘) 生效源; 默认 tickflow (需 Expert 档) */
|
||||
full_minute_data_provider?: string
|
||||
/** 分钟源 1 分钟历史深度(交易日); null/缺省 = 深历史 (如 tickflow)。分时档位据此收窄 */
|
||||
minute_history_days?: number | null
|
||||
depth5_data_provider?: string
|
||||
@@ -1814,7 +1817,7 @@ export const api = {
|
||||
}),
|
||||
}),
|
||||
|
||||
/** 全量分钟 (盘中全市场分钟落盘) 服务状态 (TickFlow Expert 专有) */
|
||||
/** 全量分钟 (盘中全市场分钟落盘) 服务状态 (按能力路由: TickFlow Expert 或声明 full_minute 的插件/自定义源) */
|
||||
minuteRefreshStatus: () =>
|
||||
request<{
|
||||
available: boolean
|
||||
@@ -1822,9 +1825,11 @@ export const api = {
|
||||
running?: boolean
|
||||
/** 读侧 freshness: 本地分区正被服务持续写入 (前端据此切本地读/解除截断) */
|
||||
healthy?: boolean
|
||||
provider?: string
|
||||
provider_effective?: string
|
||||
repair_only?: boolean
|
||||
interval_seconds?: number
|
||||
capability_ok?: boolean
|
||||
custom_provider_active?: boolean
|
||||
in_trading_hours?: boolean
|
||||
gate_reason?: string | null
|
||||
rounds?: number
|
||||
|
||||
@@ -9,7 +9,7 @@ import { toast } from '@/components/Toast'
|
||||
const INPUT_CLS =
|
||||
'w-full h-9 px-2.5 rounded-lg bg-base border-0 ring-1 ring-border/40 text-xs text-foreground placeholder:text-muted/30 focus:outline-none focus:ring-2 focus:ring-accent/40 transition-shadow'
|
||||
|
||||
const DATASETS = ['daily', 'adj_factor', 'realtime', 'minute'] as const
|
||||
const DATASETS = ['daily', 'adj_factor', 'realtime', 'minute', 'full_minute'] as const
|
||||
type DatasetKey = typeof DATASETS[number]
|
||||
|
||||
const DATASET_LABEL: Record<DatasetKey, string> = {
|
||||
@@ -17,6 +17,7 @@ const DATASET_LABEL: Record<DatasetKey, string> = {
|
||||
adj_factor: '除权因子',
|
||||
realtime: '实时行情',
|
||||
minute: '分钟K',
|
||||
full_minute: '全量分钟',
|
||||
}
|
||||
|
||||
const TARGET_FIELDS: Record<DatasetKey, string[]> = {
|
||||
@@ -24,6 +25,7 @@ const TARGET_FIELDS: Record<DatasetKey, string[]> = {
|
||||
adj_factor: ['symbol', 'trade_date', 'ex_factor'],
|
||||
realtime: ['symbol', 'name', 'last_price', 'prev_close', 'open', 'high', 'low', 'volume', 'amount', 'change_pct', 'change_amount', 'amplitude', 'turnover_rate', 'timestamp', 'session'],
|
||||
minute: ['symbol', 'datetime', 'open', 'high', 'low', 'close', 'volume', 'amount'],
|
||||
full_minute: ['symbol', 'datetime', 'open', 'high', 'low', 'close', 'volume', 'amount'],
|
||||
}
|
||||
|
||||
// 内部字段的中文说明 (下拉选项展示用)
|
||||
@@ -486,7 +488,7 @@ function DatasetDetail({
|
||||
</Field>
|
||||
</>
|
||||
)}
|
||||
{datasetKey === 'minute' && (
|
||||
{(datasetKey === 'minute' || datasetKey === 'full_minute') && (
|
||||
<>
|
||||
<Field label="资产类型参数">
|
||||
<input
|
||||
|
||||
@@ -150,6 +150,7 @@ const DEFAULT_ROUTING: Record<ProviderField, string> = {
|
||||
daily_data_provider: 'tickflow',
|
||||
adj_factor_provider: 'tickflow',
|
||||
minute_data_provider: 'tickflow',
|
||||
full_minute_data_provider: 'tickflow',
|
||||
depth5_data_provider: 'tickflow',
|
||||
realtime_data_provider: 'tickflow',
|
||||
financial_data_provider: 'tickflow',
|
||||
@@ -612,19 +613,13 @@ export function SettingsDataSourcesPanel({ highlight }: { highlight?: string } =
|
||||
daily: dailyPref,
|
||||
adj_factor: adjPref === 'same_as_daily' ? dailyPref : adjPref,
|
||||
minute: prefs.data?.minute_data_provider || 'tickflow',
|
||||
full_minute: prefs.data?.full_minute_data_provider || 'tickflow',
|
||||
realtime: prefs.data?.realtime_data_provider || 'tickflow',
|
||||
depth5: prefs.data?.depth5_data_provider || 'tickflow',
|
||||
financial: prefs.data?.financial_data_provider || 'tickflow',
|
||||
}
|
||||
const servingDatasets = (name: string) => {
|
||||
const ids = Object.entries(effProvider).filter(([, v]) => v === name).map(([k]) => k)
|
||||
if (name === 'tickflow') {
|
||||
// 不可路由能力 (field=null, 如全量分钟): 仅 TickFlow 提供, usable 即服务中
|
||||
ids.push(...(matrix.data?.capabilities ?? [])
|
||||
.filter(c => c.field == null && c.usable).map(c => c.id))
|
||||
}
|
||||
return ids
|
||||
}
|
||||
const servingDatasets = (name: string) =>
|
||||
Object.entries(effProvider).filter(([, v]) => v === name).map(([k]) => k)
|
||||
const servingSetOf = (name: string) => new Set(servingDatasets(name))
|
||||
|
||||
const matrixCaps = matrix.data?.capabilities ?? []
|
||||
|
||||
@@ -366,19 +366,19 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
|
||||
|
||||
{/* ========== 右列 ========== */}
|
||||
<div className="space-y-6">
|
||||
{/* 全量分钟 (TickFlow Expert 专有): 盘中全市场分钟落盘, intraday.universe 单请求增量 */}
|
||||
{/* 全量分钟: 盘中全市场分钟落盘, 按能力路由 (TickFlow Expert 或声明 full_minute 的插件/自定义源) */}
|
||||
<Card icon={Zap} title="全量分钟" anchor="minute-refresh">
|
||||
<ToggleRow
|
||||
label="全量分钟落盘"
|
||||
desc={
|
||||
!hasFullMinuteCap ? '需要全量分钟能力 (TickFlow Expert)'
|
||||
: rs?.custom_provider_active ? '已配置自定义分钟源, 盘中增量由插件自管'
|
||||
!hasFullMinuteCap ? '需要全量分钟能力 (TickFlow Expert 或声明该能力的自定义源)'
|
||||
: rs?.repair_only ? `服务运行中 · ${rs?.provider ?? '自定义源'} 无廉价增量端点, 按 ≥60s 全天批量节奏`
|
||||
: rs?.running ? (rs?.in_trading_hours ? '服务运行中' : '运行中 · 非连续竞价时段暂停')
|
||||
: '已关闭'
|
||||
}
|
||||
checked={prefs?.minute_refresh_enabled ?? false}
|
||||
onChange={(v) => save({ minute_refresh_enabled: v })}
|
||||
disabled={!hasFullMinuteCap || !!rs?.custom_provider_active}
|
||||
disabled={!hasFullMinuteCap}
|
||||
/>
|
||||
<div className="mt-3 pt-3 border-t border-border">
|
||||
<div className="flex items-center justify-between gap-4 py-1">
|
||||
|
||||
Reference in New Issue
Block a user