feat(minute): 全量分钟能力位与两阶段日内分钟落盘

- tickflow SDK 0.1.25: intraday.universe 标的池单请求拉全市场当日分钟
- 新能力位 Cap.INTRADAY_UNIVERSE (TickFlow Expert 专有) + 探测/别名/schema v6
- 分钟刷新两阶段: 冷启动与缺口修复走 intraday_batch 全天突发(分块容错), 稳态走 universe 增量
- 覆盖看门狗: 落后>3min / 无数据 / 连续空轮自动升级全量自愈
- 刷新间隔钳制 [3,300]s 默认 6s; 监控页全量分钟开关与状态入口
- CONTRIBUTING: 分钟 K 北京时间墙钟契约 (naive, 入口强制归一)
This commit is contained in:
shy3130
2026-08-30 19:05:25 +08:00
parent 246a7df57d
commit 657d4d9948
18 changed files with 817 additions and 183 deletions
+1
View File
@@ -106,6 +106,7 @@
- 窗口、前 N 日和批次回算均按实际交易日,不得用自然日直接替代。
- A 股交易时段统一按北京时间处理;服务器时区不能成为业务逻辑的隐式输入。
- 分钟 K 的 `datetime` 统一为北京时间墙钟(naive,如 `09:35:00`);数据源入口(`kline_sync``_normalize_minute` / `_try_custom_minute`)强制归一,禁止 UTC 口径入库或下发。
- 日线、分钟线和实时快照必须明确交易日期归属,尤其注意午休、收盘后和跨日重启。
- 分钟 K 的股票、ETF、指数分开存储和路由,不得仅凭代码格式猜测资产类型。
+4 -2
View File
@@ -1095,9 +1095,11 @@ async def clear_minute(request: Request):
removed = 0
if minute_dir.exists():
try:
result = repo.db.execute("SELECT COUNT(*) AS cnt FROM kline_minute").fetchone()
# execute_one (cursor+close): 直连 db.execute 的未消费结果集会在 Windows 上
# 钉住分区句柄, 导致下方 rmtree 静默删不掉被钉文件
result = repo.execute_one("SELECT COUNT(*) AS cnt FROM kline_minute")
removed = result[0] if result else 0
except Exception: # noqa: BLE001
except Exception:
pass
# 仅删 kline_minute 目录, 绝不触碰其他目录
shutil.rmtree(minute_dir, ignore_errors=True)
+18 -19
View File
@@ -386,6 +386,21 @@ def _realtime_allowed() -> bool:
return QuoteService.is_realtime_allowed()
def _minute_history_days() -> int | None:
"""当前分钟源的 1 分钟历史深度(交易日); None = 深历史(tickflow 基准)。
provider 可选类属性 minute_history_days 声明 (如 stock-sdk = 5,
免费分时接口只保留最近 5 个交易日); 未声明或走 tickflow 时视为深历史。
前端分时档位/默认值据此收窄。
"""
from app.services import kline_sync, preferences
provider_name = preferences.get_minute_data_provider()
provider, fallback, _err = kline_sync._resolve_minute_provider(provider_name)
if fallback or provider is None:
return None
return getattr(provider, "minute_history_days", None)
class MinuteSyncPrefs(BaseModel):
minute_sync_enabled: bool
minute_sync_days: int = 5
@@ -484,12 +499,12 @@ 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(),
"minute_history_days": _minute_history_days(),
"depth5_data_provider": preferences.get_depth5_data_provider(),
"realtime_data_provider": preferences.get_realtime_data_provider(),
"financial_data_provider": preferences.get_financial_provider(),
"data_source_job_timeout_s": preferences.get_data_source_job_timeout_s(),
"data_source_long_job_timeout_s": preferences.get_data_source_long_job_timeout_s(),
"realtime_watchlist_symbols": preferences.get_realtime_watchlist_symbols(),
**preferences.get_realtime_quote_scope(),
"pipeline_pull_a_share": preferences.get_pipeline_pull_a_share(),
"pipeline_pull_etf": preferences.get_pipeline_pull_etf(),
@@ -895,8 +910,8 @@ class RealtimeQuoteScopePrefs(BaseModel):
def update_realtime_quotes(req: RealtimeQuotesPrefs, request: Request) -> dict:
"""保存全局实时行情开关。
none 档无实时行情权限;free 档开启自选股实时;starter+ 开启全市场实时。
前端据此把开关置灰 / 回弹。
无实时能力的档位(TickFlow none/free)开关回弹强制关闭;
starter+ 或自定义实时源(如 fuyao)为全市场实时。前端据此把开关置灰 / 回弹。
"""
from app.services import preferences
qs = getattr(request.app.state, "quote_service", None)
@@ -950,10 +965,6 @@ def update_realtime_quotes(req: RealtimeQuotesPrefs, request: Request) -> dict:
+ f"(任务 {job_id}"
)
raise HTTPException(status_code=409, detail=detail)
if req.realtime_quotes_enabled and qs and qs.realtime_mode() == "watchlist" and not preferences.get_realtime_watchlist_symbols():
preferences.save({"realtime_quotes_enabled": False})
_sync_depth_polling(False)
return {"realtime_quotes_enabled": False, "realtime_allowed": True, "mode": "watchlist", "error": "watchlist_empty"}
preferences.save({"realtime_quotes_enabled": req.realtime_quotes_enabled})
if qs:
@@ -974,18 +985,6 @@ def update_realtime_quote_scope(req: RealtimeQuoteScopePrefs) -> dict:
return preferences.set_realtime_quote_scope(cfg)
class RealtimeWatchlistPrefs(BaseModel):
symbols: list[str] = []
@router.put("/preferences/realtime-watchlist")
def update_realtime_watchlist(req: RealtimeWatchlistPrefs) -> dict:
"""兼容旧入口;Free 实时标的由自选页前 5 个决定。"""
from app.services import preferences
symbols = preferences.set_realtime_watchlist_symbols(req.symbols)
return {"realtime_watchlist_symbols": symbols}
class IndicesNavPinnedPrefs(BaseModel):
indices_nav_pinned: bool
+162 -14
View File
@@ -20,7 +20,7 @@ from app.services import preferences
from app.tickflow.capabilities import Cap, CapabilitySet
from app.tickflow.client import get_client
from app.tickflow.rate_limits import chunked, resolve_limit, sleep_between_batches
from app.tickflow.repository import KlineRepository
from app.tickflow.repository import KlineRepository, replace_with_retry
logger = logging.getLogger(__name__)
@@ -32,11 +32,12 @@ def _atomic_write_parquet(df: pl.DataFrame, out) -> None:
单文件、每次「读→concat→原地写」, 直接 write_parquet(out) 在进程被 kill
(dev.sh 清端口用 kill -9)、reap 超时或断电时会留下半截文件, 之后复权视图
scan_parquet 整条链路报错、enriched 全市场重算不出。临时文件后缀 .tmp 不匹配
*.parquet glob, 不会被扫描误读。
*.parquet glob, 不会被扫描误读。Windows 下目标正被并发读取时由
replace_with_retry 短退避穿过。
"""
tmp = out.with_name(out.name + ".tmp")
df.write_parquet(tmp)
tmp.replace(out) # 同目录 rename, POSIX/NTFS 均为原子操作
replace_with_retry(tmp, out)
# 标准列(无论 SDK 返回什么形状,我们把它规范成这套)
@@ -444,8 +445,80 @@ CANONICAL_MINUTE_COLS = [
]
# 北京墙钟特征时段(含集合竞价 09:15 与收盘 15:00): 上午 09-11, 下午 13-15
_BJ_HOURS = [9, 10, 11, 13, 14, 15]
# 上述时段 -8h 的 UTC 墙钟特征: 上午 01-03, 下午 05-07
_UTC_SHIFTED_HOURS = [1, 2, 3, 5, 6, 7]
def _enforce_minute_beijing_wallclock(df: pl.DataFrame, *, source: str) -> pl.DataFrame:
"""分钟 K datetime 时区契约守卫: 统一为北京墙钟 (naive)。
契约 (CONTRIBUTING §3.3): kline_minute.datetime 必须是北京时间墙钟, 如 09:35:00。
在两个源头入口强制 —— _normalize_minute (TickFlow 帧) 与 _try_custom_minute
(插件/自定义源帧); 落盘 (_write_minute_partition) 与内存消费 (监控/补拉/脉冲)
均在其下游, 这里收口即全覆盖:
- tz-aware → 转 Asia/Shanghai 后去时区;
- naive 且时刻落在 A 股交易时段 → 直通 (已是北京墙钟);
- naive 且整体呈"交易时段 -8h"的 UTC 特征 → 自动 +8 纠偏并记日志;
- 无法识别的口径 → fail-closed 抛 ValueError, 不让脏时间入库或下发。
幂等: 纠偏后的帧再过守卫直通, 不会二次改写。
"""
if df.is_empty() or "datetime" not in df.columns:
return df
dtype = df.schema["datetime"]
if not isinstance(dtype, pl.Datetime):
# trade_time 等字符串路径: 先解析成 Datetime (失败置 null), 再做时段分类
if dtype == pl.Utf8:
df = df.with_columns(pl.col("datetime").str.to_datetime(strict=False))
else:
df = df.with_columns(pl.col("datetime").cast(pl.Datetime("us"), strict=False))
dtype = df.schema["datetime"]
if not isinstance(dtype, pl.Datetime):
return df # 仍非 Datetime: 维持原行为交由下游处理
if isinstance(dtype, pl.Datetime) and dtype.time_zone is not None:
df = df.with_columns(
pl.col("datetime")
.dt.convert_time_zone("Asia/Shanghai")
.dt.replace_time_zone(None)
.cast(pl.Datetime("us"))
)
logger.info("minute datetime tz-aware input converted to Beijing wallclock (source=%s)", source)
return df
hour = pl.col("datetime").dt.hour()
beijing = int(df.select(hour.is_in(_BJ_HOURS).sum()).item() or 0)
utc_shifted = int(df.select(hour.is_in(_UTC_SHIFTED_HOURS).sum()).item() or 0)
if beijing == 0 and utc_shifted == 0:
if df["datetime"].null_count() == df.height:
return df # 全 null: 维持原行为, 由下游落盘过滤
raise ValueError(
f"minute datetime 口径无法识别 (source={source}, rows={df.height}, "
f"sample={df['datetime'].drop_nulls().head(2).to_list()}): "
"契约要求北京墙钟 (09:30-15:00), 既非交易时段也非 UTC 平移特征"
)
if utc_shifted > beijing:
if beijing:
logger.warning(
"minute datetime mixed convention, shifting all by +8h per UTC majority "
"(source=%s, utc=%d, beijing=%d)", source, utc_shifted, beijing,
)
else:
logger.info(
"minute datetime UTC wallclock detected, shifted +8h to Beijing "
"(source=%s, rows=%d)", source, utc_shifted,
)
return df.with_columns(pl.col("datetime") + pl.duration(hours=8))
if utc_shifted:
logger.warning(
"minute datetime has %d UTC-like rows among %d Beijing rows, left as-is "
"(source=%s)", utc_shifted, beijing, source,
)
return df
def _normalize_minute(df_in, default_symbol: str | None = None) -> pl.DataFrame:
"""把 SDK 返回的分钟 K 数据规范成 canonical 列。"""
"""把 SDK 返回的分钟 K 数据规范成 canonical 列 (datetime 收口为北京墙钟)"""
if df_in is None or len(df_in) == 0:
return pl.DataFrame()
@@ -463,8 +536,14 @@ def _normalize_minute(df_in, default_symbol: str | None = None) -> pl.DataFrame:
# datetime 列:优先用 timestamp(毫秒精度),其次 trade_time
if "timestamp" in df.columns:
# TickFlow 毫秒时间戳为 UTC 基准; 契约要求北京墙钟 naive
# (与 stock-sdk provider 归一口径一致, 见 CONTRIBUTING §3.3)
df = df.with_columns(
pl.from_epoch("timestamp", time_unit="ms").alias("datetime"),
pl.from_epoch(pl.col("timestamp").cast(pl.Int64), time_unit="ms")
.dt.replace_time_zone("UTC")
.dt.convert_time_zone("Asia/Shanghai")
.dt.replace_time_zone(None)
.alias("datetime")
).drop("timestamp")
for drop_col in ("trade_time", "trade_date"):
if drop_col in df.columns:
@@ -476,6 +555,10 @@ def _normalize_minute(df_in, default_symbol: str | None = None) -> pl.DataFrame:
elif "trade_date" in df.columns:
df = df.rename({"trade_date": "datetime"})
if "datetime" in df.columns:
# 时区契约守卫: 须在下方通用 us-cast 之前, 避免带时区列被静默剥成 UTC-naive
df = _enforce_minute_beijing_wallclock(df, source="tickflow")
if "symbol" not in df.columns and default_symbol is not None:
df = df.with_columns(pl.lit(default_symbol).alias("symbol"))
@@ -604,11 +687,18 @@ def _try_custom_minute(
symbols, start_time=start_time, end_time=end_time,
asset_type=asset_type, freq=freq, on_chunk_done=wrapped_cb,
)
return (df, False)
except Exception as e: # noqa: BLE001
except Exception as e:
logger.warning("custom minute provider %s call failed, falling back to TickFlow: %s",
provider_name, e)
return (None, True)
try:
# 时区契约守卫: 插件/自定义源帧同样收口为北京墙钟 (CONTRIBUTING §3.3)
df = _enforce_minute_beijing_wallclock(df, source=provider_name)
except Exception as e:
logger.warning("custom minute provider %s datetime 契约校验失败, falling back to TickFlow: %s",
provider_name, e)
return (None, True)
return (df, False)
def sync_minute_batch(
@@ -838,7 +928,10 @@ def fetch_intraday_full_market_burst(
- 监控路径每轮只拉少量标的 (≤ batch 上限, 单请求);
本函数按 batch_size 把全市场切块后用线程池一次全部打出
(5546/200 = 28 并发), 配合 >=60s 的固定轮节奏, 任何 60s
滑动窗口至多一个脉冲 (28 < 48 安全 rpm), 轮内失败不重试
滑动窗口至多一个脉冲 (28 < 48 安全 rpm)。
- 单块失败不拖垮整轮: 成功块照常返回落盘; 失败块立即单独重试一次,
仍失败则跳过该块 (本函数每轮拉全天, 下一轮天然自愈)。失败块过多
(>4, 系统性故障/限流风暴) 时跳过重试, 避免向已过载的服务端加压。
限流口径: 只用 intraday.batch 独立池 (Cap.INTRADAY_BATCH, Expert 专有),
不与 kline.minute.batch (盘后分钟同步) 共享配额。
@@ -856,20 +949,75 @@ def fetch_intraday_full_market_burst(
tf = get_client()
def _fetch(chunk: list[str]) -> list[pl.DataFrame]:
def _fetch(chunk: list[str]) -> tuple[list[pl.DataFrame], Exception | None]:
# 单块独立容错: 异常作为返回值上交而不是抛出, 避免一个块把整轮
# 已成功的数据一起拖垮 (pool.map 迭代中抛异常会废弃全部已收 frames)
try:
raw = tf.klines.intraday_batch(
chunk, count=count, as_dataframe=True, show_progress=False,
batch_size=len(chunk),
)
return _normalize_intraday_raw(raw)
return (_normalize_intraday_raw(raw), None)
except Exception as e:
return ([], e)
frames: list[pl.DataFrame] = []
failed: list[list[str]] = []
with ThreadPoolExecutor(max_workers=min(len(chunks), 32)) as pool:
for result in pool.map(_fetch, chunks):
frames.extend(result)
for chunk, (sub, err) in zip(chunks, pool.map(_fetch, chunks), strict=True):
if err is not None:
failed.append(chunk)
else:
frames.extend(sub)
requests = len(chunks)
if failed:
if len(failed) > 4:
missed = sum(len(chunk) for chunk in failed)
logger.warning(
"intraday burst: %d/%d chunks failed (%d symbols), systemic — skip retry, next round re-pulls full day",
len(failed), len(chunks), missed,
)
else:
logger.warning("intraday burst: %d/%d chunks failed, retrying once", len(failed), len(chunks))
for chunk in failed:
sub, err = _fetch(chunk)
requests += 1
if err is None:
frames.extend(sub)
else:
logger.warning(
"intraday burst: chunk retry still failed, skip %d symbols this round: %s",
len(chunk), err,
)
if not frames:
return (pl.DataFrame(), len(chunks))
return (pl.concat(frames, how="diagonal_relaxed"), len(chunks))
return (pl.DataFrame(), requests)
return (pl.concat(frames, how="diagonal_relaxed"), requests)
def fetch_intraday_universe_increment(
universe: str = "CN_Equity_A",
*,
count: int = 3,
) -> tuple[pl.DataFrame, int]:
"""全市场当日分钟K增量拉取 (盘中稳态轮专用, 不落盘)。
/v1/klines/intraday/universe: 传 universe ID 一次请求返回全市场每只标的
最新 count 根分钟K (服务端实测上限 3 根/标的), 替代稳态场景下 28 块并发
的 intraday.batch 脉冲 (请求量 28→1, 传输量 ~40 倍降)。缺口回补
(冷启动/长时间断档/全天修复) 仍走 fetch_intraday_full_market_burst。
返回 (增量分钟K, 请求数); 拉取失败返回空 df 由调用方按失败轮处理。
"""
tf = get_client()
try:
raw = tf.klines.intraday_universe(universe, count=count, as_dataframe=True)
except Exception as e:
logger.warning("intraday universe fetch failed (%s): %s", universe, e)
return (pl.DataFrame(), 0)
frames = _normalize_intraday_raw(raw)
if not frames:
return (pl.DataFrame(), 0)
return (pl.concat(frames, how="diagonal_relaxed"), 1)
def fetch_minute_single(
+87 -21
View File
@@ -1,26 +1,39 @@
"""盘中分钟K增量落盘服务 (Expert 专有)。
每轮用 intraday.batch (日内分时批量, 独立限流池) 并发脉冲拉全市场当日分钟K,
单次合并写入当日 kline_minute 分区, 供分钟策略 (minute_filter) 读到新鲜数据。
两段式拉取 (见 feat/minute-strategy 方案):
- 全天修复轮: intraday.batch (日内分时批量) 并发脉冲一次拉全市场当日全部
分钟K — 冷启动 (如 10 点才开服务, 补 9:30 起缺口) / 覆盖滞后超阈值 /
连续空轮自愈时触发。
- 稳态增量轮: intraday.universe 传 CN_Equity_A 标的池, 单请求返回全市场
每只最新 3 根 (服务端上限), 靠 _write_minute_partition 的
unique(symbol,datetime) 幂等合并滚出全天。
设计约束 (见 feat/minute-strategy 方案):
- Expert 专有: 能力门控 Cap.INTRADAY_BATCH — 该能力仅 Expert 档具备, 天然排他。
- 并发脉冲: 全市场按 batch_size 分块 (5546/200 = 28 块), ThreadPoolExecutor 一次
打出全部块 (≤28 并发)。任何 60s 滑动窗口至多一个脉冲 (28 < 48 安全 rpm)。
- 固定节奏: 默认 60s 一轮 (clamp [60, 300]), 下一轮 = max(本轮起点+间隔, 上轮完成),
不补跑 (missed 轮次直接跳过), 轮内失败不重试
- 仅连续竞价时段运行 (9:30-11:30 / 13:00-15:00), 午休/收盘自动暂停与恢复。
单轮合并写入当日 kline_minute 分区, 供分钟策略 (minute_filter) 读到新鲜数据。
设计约束:
- Expert 专有: 能力门控 Cap.INTRADAY_UNIVERSE (全量分钟) — 仅 TickFlow
Expert 档具备, 天然排他 (自定义分钟源无此能力, 且服务本就让位插件)。
- 修复轮并发脉冲: 全市场按 batch_size 分块 (5546/200 = 28 块) 一次打出
任何 60s 滑动窗口至多一个脉冲 (28 < 48 安全 rpm); 单块失败不拖垮整轮,
失败块单独重试一次 (见 fetch_intraday_full_market_burst)。
- 稳态轮单请求: 无脉冲并发, 间隔可低至 3s; 实际节奏 = max(间隔, 单轮完成),
服务端响应 ~5s 时自动退化为响应节奏, 不会重叠请求。
- 固定节奏: 默认 6s 一轮 (clamp [3, 300]), 不补跑 (missed 轮次直接跳过)。
- 仅连续竞价时段运行 (9:30-11:30 / 13:00-15:00), 午休/收盘自动暂停与恢复;
午休后恢复因覆盖滞后会多跑一次修复轮, 幂等无害。
- 不与其他分钟能力冲突: 与 盘后分钟同步 (kline.minute.batch) / 分时监控路径
(fetch_intraday_monitor_batch) 分属不同限流池; 落盘走 _write_minute_partition
的 unique(symbol,datetime) 合并, 与盘后同步写同一分区安全幂等。
分属不同限流池; 落盘走 _write_minute_partition 的 unique(symbol,datetime)
合并, 与盘后同步写同一分区安全幂等。
- 数据源插件化让位: 配置了自定义分钟源 (minute_data_provider != tickflow) 时
服务不启动 — 盘中增量交由插件自管, 本服务不抢占。
分层: 本模块只做调度/落盘/状态; TickFlow SDK 调用全部在 kline_sync 边界层
(fetch_intraday_full_market_burst), 保持插件化边界不泄漏。
(fetch_intraday_full_market_burst / fetch_intraday_universe_increment),
保持插件化边界不泄漏。
"""
from __future__ import annotations
import contextlib
import threading
import time
from dataclasses import dataclass, field
@@ -28,14 +41,18 @@ from typing import Any
import polars as pl
from app.market_time import in_continuous_session
from app.market_time import cn_now, cn_today, in_continuous_session
from app.services import preferences
# 轮询间隔允许范围 (秒): 下限 60s 保证任何滑动窗口 ≤1 个脉冲, 上限防误配。
REFRESH_INTERVAL_MIN = 60
# 轮询间隔允许范围 (秒): 稳态轮单请求无并发脉冲, 下限 3s; 上限防误配。
REFRESH_INTERVAL_MIN = 3
REFRESH_INTERVAL_MAX = 300
# 等待步长 (秒): 循环小步睡眠, 便于快速停止与偏好热生效。
_LOOP_STEP_S = 2.0
# 当日覆盖滞后超过该分钟数 (≈ universe 单请求 3 根余量) → 触发全天修复轮。
_REPAIR_LAG_MINUTES = 3.0
# 连续空轮达到该次数 → 强制全天修复轮 (自愈 universe 端点持续异常)。
_EMPTY_ROUNDS_TO_REPAIR = 2
def _in_continuous_session(now=None) -> bool:
@@ -52,7 +69,8 @@ class _RefreshState:
last_round_ms: float | None = None # 单轮耗时
last_rows: int = 0 # 上轮写入行数 (合并后)
last_symbols: int = 0 # 上轮覆盖标的数
last_requests: int = 0 # 上轮请求数 (分块数)
last_requests: int = 0 # 上轮请求数 (增量恒 1, 修复=分块数+重试)
last_mode: str | None = None # 上轮模式: "increment" / "full"
last_error: str | None = None
next_round_at: float | None = None # epoch 秒
extra: dict[str, Any] = field(default_factory=dict)
@@ -68,6 +86,7 @@ class MinuteRefreshService:
self._stop = threading.Event()
self._state = _RefreshState()
self._round_lock = threading.Lock() # 同时只允许一轮 (手动触发与定时轮互斥)
self._empty_rounds = 0 # 连续空轮计数 (escalate 到全天修复)
# ------------------------------------------------------------------
# 生命周期
@@ -98,14 +117,18 @@ class MinuteRefreshService:
# ------------------------------------------------------------------
def capability_ok(self) -> bool:
"""Cap.INTRADAY_BATCH 存在 (Expert)。能力探测结果缓存在 app.state。"""
"""Cap.INTRADAY_UNIVERSE (全量分钟) 存在。能力探测结果缓存在 app.state。
门控挂在稳态增量的主能力上; 全天修复轮用的 intraday.batch 与其
同属 Expert 档 (tiers.yaml), 目前两者必然同时持有。
"""
capset = getattr(self._app_state, "capabilities", None) if self._app_state else None
if capset is None:
return False
try:
from app.tickflow.capabilities import Cap
return capset.has(Cap.INTRADAY_BATCH)
return capset.has(Cap.INTRADAY_UNIVERSE)
except Exception:
return False
@@ -162,23 +185,64 @@ class MinuteRefreshService:
# 单轮
# ------------------------------------------------------------------
def _today_coverage_lag_minutes(self) -> float | None:
"""当日分区最新K距现在的分钟数; None = 当日无数据。
覆盖度探测只读当日分区文件 (毫秒级); 任何异常按无数据处理 →
本轮走全天修复, 不会因探测失败而丢增量。
"""
with contextlib.suppress(Exception):
part = (
self._repo.store.data_dir / "kline_minute"
/ f"date={cn_today().isoformat()}" / "part.parquet"
)
if not part.exists():
return None
mx = pl.read_parquet(part, columns=["datetime"])["datetime"].max()
if mx is None:
return None
# 分区 datetime 为北京墙钟 naive, cn_now 带时区 → 剥齐再比
return (cn_now().replace(tzinfo=None) - mx).total_seconds() / 60
return None
def _select_mode(self) -> str:
"""选轮次模式: 稳态增量 (universe 单请求) vs 全天修复 (burst 脉冲)。
当日已有数据且覆盖滞后 ≤ _REPAIR_LAG_MINUTES (≈ 3 根余量) → 增量;
冷启动 / 断档超阈值 / 连续空轮 → 全天修复。
"""
if self._empty_rounds >= _EMPTY_ROUNDS_TO_REPAIR:
return "full"
lag = self._today_coverage_lag_minutes()
if lag is None or lag > _REPAIR_LAG_MINUTES:
return "full"
return "increment"
def _run_round(self) -> None:
from app.services import kline_sync
t0 = time.perf_counter()
mode = self._select_mode()
with self._round_lock:
if mode == "increment":
df, requests = kline_sync.fetch_intraday_universe_increment()
self._state.last_symbols = (
df["symbol"].n_unique() if not df.is_empty() else 0
)
else:
symbols = self._universe()
self._state.last_symbols = len(symbols)
if not symbols:
self._state.last_error = "empty universe (instruments 未加载)"
return
capset = getattr(self._app_state, "capabilities", None) if self._app_state else None
with self._round_lock:
df, requests = kline_sync.fetch_intraday_full_market_burst(symbols, capset)
self._state.last_requests = requests
if df.is_empty():
self._state.last_error = "intraday burst returned no data"
self._empty_rounds += 1
self._state.last_error = f"intraday {mode} returned no data"
return
self._empty_rounds = 0
written = kline_sync._write_minute_partition(
df, self._repo.store.data_dir / "kline_minute",
)
@@ -187,6 +251,7 @@ class MinuteRefreshService:
self._state.last_round_at = time.time()
self._state.last_round_ms = (time.perf_counter() - t0) * 1000
self._state.last_rows = written
self._state.last_mode = mode
self._state.last_error = None
def _universe(self) -> list[str]:
@@ -221,6 +286,7 @@ class MinuteRefreshService:
"last_rows": self._state.last_rows,
"last_symbols": self._state.last_symbols,
"last_requests": self._state.last_requests,
"last_mode": self._state.last_mode,
"next_round_at": self._state.next_round_at,
"last_error": self._state.last_error,
}
+8 -29
View File
@@ -84,29 +84,6 @@ def get_realtime_quote_interval() -> float:
return load().get("realtime_quote_interval", 6.0)
def get_realtime_watchlist_symbols() -> list[str]:
"""Free 档自选实时监控标的:直接取自选页前 5 个。"""
try:
from app.services import watchlist
rows = watchlist.list_symbols()
except Exception as e: # noqa: BLE001
logger.warning("load watchlist for realtime failed: %s", e)
return []
out: list[str] = []
for row in rows:
symbol = str((row or {}).get("symbol") or "").strip().upper()
if symbol and symbol not in out:
out.append(symbol)
if len(out) >= 5:
break
return out
def set_realtime_watchlist_symbols(symbols: list[str]) -> list[str]: # noqa: ARG001
"""兼容旧接口: Free 实时标的现在由自选页前 5 个决定。"""
return get_realtime_watchlist_symbols()
def set_realtime_quote_interval(interval: float) -> float:
"""保存行情轮询间隔(不在此做 min/max 校验,由调用方按档位限制)。"""
current = load()
@@ -214,10 +191,12 @@ def get_minute_sync_segment_days() -> int:
"""
return max(5, min(30, load().get("minute_sync_segment_days", 20)))
# ===== 盘中分钟增量刷新 (Expert 专有, intraday.batch 独立限流池) =====
# ===== 盘中分钟增量刷新 (Expert 专有) =====
# 下限 60s: 保证任何 60s 滑动窗口至多一个全市场脉冲 (28 并发 < 48 安全 rpm)。
_MINUTE_REFRESH_INTERVAL_MIN = 60
# 稳态轮为 intraday.universe 单请求增量, 无脉冲并发, 间隔可低至 3s;
# 全天修复轮 (intraday.batch 28 块爆发) 的 rpm 安全与间隔无关, 由轮次
# 调度 max(间隔, 单轮完成) 天然防重叠。
_MINUTE_REFRESH_INTERVAL_MIN = 3
_MINUTE_REFRESH_INTERVAL_MAX = 300
@@ -227,10 +206,10 @@ def get_minute_refresh_enabled() -> bool:
def get_minute_refresh_interval() -> int:
"""盘中分钟增量刷新间隔(秒)。默认 60,范围 [60, 300]。"""
"""盘中分钟增量刷新间隔(秒)。默认 6,范围 [3, 300]。"""
return max(
_MINUTE_REFRESH_INTERVAL_MIN,
min(_MINUTE_REFRESH_INTERVAL_MAX, int(load().get("minute_refresh_interval", 60))),
min(_MINUTE_REFRESH_INTERVAL_MAX, int(load().get("minute_refresh_interval", 6))),
)
@@ -947,7 +926,7 @@ def set_realtime_monitor_config(cfg: dict) -> dict:
if "minute_refresh_enabled" in cfg:
updates["minute_refresh_enabled"] = bool(cfg["minute_refresh_enabled"])
if "minute_refresh_interval" in cfg:
# clamp 到 [60, 300] (下限保证 60s 窗口至多一个全市场脉冲), 与 getter 一致
# clamp 到 [3, 300], 与 getter 一致, 防前端传越界值
updates["minute_refresh_interval"] = max(
_MINUTE_REFRESH_INTERVAL_MIN,
min(_MINUTE_REFRESH_INTERVAL_MAX, int(cfg["minute_refresh_interval"])))
+1
View File
@@ -20,6 +20,7 @@ class Cap(StrEnum):
KLINE_MINUTE_BATCH = "kline.minute.batch"
INTRADAY = "intraday"
INTRADAY_BATCH = "intraday.batch"
INTRADAY_UNIVERSE = "intraday.universe"
DEPTH5 = "depth5"
DEPTH5_BATCH = "depth5.batch"
WEBSOCKET = "websocket"
+8 -1
View File
@@ -32,7 +32,8 @@ _CAPSET_CACHE_FILE = "capabilities.json"
# v2: 拆分 depth5 → depth5(单只) + depth5.batch(批量)
# v3: 探测补全 quote.batch(此前 tiers.yaml 声明了但 _probe_real 漏探测)
# v5: Free 档补充付费服务器 quote.by_symbol(10rpm/5标的),用于自选股实时监控。
_CACHE_SCHEMA_VERSION = 5
# v6: 新增 intraday.universe(全量分钟) 探测。
_CACHE_SCHEMA_VERSION = 6
# 探测用最小代价请求:挑流通性最好的 1 只标的试
_PROBE_SYMBOL = "600000.SH" # 浦发银行,长期不会退市
@@ -234,6 +235,11 @@ def _probe_real(tiers: dict) -> tuple[CapabilitySet, list[str], set[Cap]]:
lambda: tf.klines.intraday_batch([_PROBE_SYMBOL], count=1, as_dataframe=False),
defaults(Cap.INTRADAY_BATCH))
# intraday.universe — 全量分钟: 标的池单请求拉全市场最新 N 根 (Expert)
try_call(Cap.INTRADAY_UNIVERSE,
lambda: tf.klines.intraday_universe("CN_Equity_A", count=1, as_dataframe=False),
defaults(Cap.INTRADAY_UNIVERSE))
# depth5 — 按标的查(单只)
try_call(Cap.DEPTH5,
lambda: tf.depth.get(_PROBE_SYMBOL),
@@ -464,6 +470,7 @@ _CAP_ALIASES: dict[Cap, str] = {
Cap.KLINE_MINUTE_BY_SYMBOL: "分钟K",
Cap.INTRADAY: "分时",
Cap.INTRADAY_BATCH: "批量分时",
Cap.INTRADAY_UNIVERSE: "全量分钟",
Cap.DEPTH5: "五档",
Cap.DEPTH5_BATCH: "批量五档",
Cap.WEBSOCKET: "WS",
@@ -0,0 +1,52 @@
"""全量分钟能力 (Cap.INTRADAY_UNIVERSE) 契约。
- 能力位存在且值为 "intraday.universe"
- tiers.yaml 仅 expert 档声明该能力 (Pro/自定义源天然没有)
- 探测层注册了该能力的探测调用, 显示标签为「全量分钟」
- 缓存 schema 已 bump (旧 capabilities.json 触发重探测)
- 盘中分钟服务门控挂在该能力位上
"""
from pathlib import Path
from types import SimpleNamespace
import polars as pl
from app.services.minute_refresh import MinuteRefreshService
from app.tickflow.capabilities import Cap, CapabilityLimits, CapabilitySet
from app.tickflow.policy import _CACHE_SCHEMA_VERSION, _CAP_ALIASES, _load_tiers_yaml
def test_capability_enum_value():
assert Cap("intraday.universe") is Cap.INTRADAY_UNIVERSE
def test_tiers_yaml_grants_universe_to_expert_only():
tiers = _load_tiers_yaml()
assert "intraday.universe" in tiers["expert"]
for tier in ("free", "starter", "pro"):
assert "intraday.universe" not in tiers[tier]
def test_policy_labels_and_cache_schema():
assert _CAP_ALIASES[Cap.INTRADAY_UNIVERSE] == "全量分钟"
assert _CACHE_SCHEMA_VERSION >= 6
def test_service_gate_requires_universe_not_batch_alone():
class _Repo:
store = SimpleNamespace(data_dir=Path("."))
def get_instruments(self) -> pl.DataFrame:
return pl.DataFrame({"symbol": []})
def _svc_with(caps: dict) -> MinuteRefreshService:
svc = MinuteRefreshService(_Repo())
svc.set_app_state(SimpleNamespace(capabilities=CapabilitySet(caps)))
return svc
universe = {Cap.INTRADAY_UNIVERSE: CapabilityLimits(rpm=20)}
batch_only = {Cap.INTRADAY_BATCH: CapabilityLimits(rpm=60, batch=200)}
assert _svc_with(universe).capability_ok() is True
assert _svc_with(batch_only).capability_ok() is False # 仅有 intraday.batch 不放行
assert _svc_with({}).capability_ok() is False
@@ -0,0 +1,90 @@
"""fetch_intraday_full_market_burst 单块容错契约。
一个块失败不得拖垮整轮: 成功块必须照常返回供落盘; 失败块单独重试一次;
失败块过多 (系统性故障) 时跳过重试。全部用假 client, 不发真实网络请求。
"""
from types import SimpleNamespace
from unittest.mock import patch
import polars as pl
from app.services import kline_sync
from app.tickflow.capabilities import Cap, CapabilityLimits, CapabilitySet
def _capset(batch: int = 2) -> CapabilitySet:
return CapabilitySet({Cap.INTRADAY_BATCH: CapabilityLimits(rpm=60, batch=batch)})
def _frame() -> pl.DataFrame:
# _normalize_minute 的最小输入: 毫秒 timestamp → 北京墙钟 datetime。
# 时间必须落在交易时段 (时区契约守卫会拒绝非交易小时的脏数据)
from datetime import datetime
from zoneinfo import ZoneInfo
ts = int(datetime(2026, 8, 28, 9, 31, tzinfo=ZoneInfo("Asia/Shanghai")).timestamp() * 1000)
return pl.DataFrame({
"timestamp": [ts],
"open": [1.0], "high": [1.0], "low": [1.0], "close": [1.0],
"volume": [100.0], "amount": [100.0],
})
class _FakeKlines:
"""intraday_batch 假实现: fail_once 首次失败重试成功, fail_always 恒失败。"""
def __init__(self, fail_once: set[str], fail_always: set[str]) -> None:
self.fail_once = set(fail_once)
self.fail_always = set(fail_always)
self.calls: list[list[str]] = []
def intraday_batch(self, chunk, **kwargs):
self.calls.append(list(chunk))
syms = set(chunk)
if syms & self.fail_always:
raise RuntimeError("permanent failure")
if syms & self.fail_once:
self.fail_once -= syms
raise RuntimeError("transient failure")
return {s: _frame() for s in chunk}
def _run(symbols: list[str], fake: _FakeKlines, batch: int = 2):
client = SimpleNamespace(klines=fake)
with patch.object(kline_sync, "get_client", return_value=client):
return kline_sync.fetch_intraday_full_market_burst(symbols, _capset(batch))
def test_transient_chunk_failure_retried_and_all_symbols_returned():
symbols = [f"S{i}" for i in range(6)] # 3 chunks (batch=2)
fake = _FakeKlines(fail_once={"S0", "S1"}, fail_always=set())
df, requests = _run(symbols, fake)
# 失败块 (S0,S1) 重试后成功 → 6 只全在, 请求数 = 3 块 + 1 次重试
assert set(df["symbol"].to_list()) == set(symbols)
assert requests == 4
def test_permanent_chunk_failure_skipped_without_losing_other_chunks():
symbols = [f"S{i}" for i in range(6)]
fake = _FakeKlines(fail_once=set(), fail_always={"S4", "S5"})
df, requests = _run(symbols, fake)
# 失败块重试仍失败 → 只跳过该块, 其余 4 只必须返回 (旧实现会整轮丢弃)
assert set(df["symbol"].to_list()) == {"S0", "S1", "S2", "S3"}
assert requests == 4
def test_all_chunks_succeed_requests_equals_chunk_count():
symbols = [f"S{i}" for i in range(6)]
fake = _FakeKlines(fail_once=set(), fail_always=set())
df, requests = _run(symbols, fake)
assert set(df["symbol"].to_list()) == set(symbols)
assert requests == 3
def test_systemic_failure_skips_retry_to_avoid_pressuring_overloaded_server():
symbols = [f"S{i}" for i in range(12)] # 6 chunks (batch=2)
# 5 个块恒失败 (>4) → 系统性故障, 不再重试
fake = _FakeKlines(fail_once=set(), fail_always={s for s in symbols if s not in ("S0", "S1")})
df, requests = _run(symbols, fake)
assert set(df["symbol"].to_list()) == {"S0", "S1"}
assert requests == 6 # 无重试
assert len(fake.calls) == 6
+62
View File
@@ -0,0 +1,62 @@
"""分钟源历史深度能力 (minute_history_days) 契约测试。
provider 可选类属性 minute_history_days 声明 1 分钟历史深度(交易日):
- stock-sdk = 5 (免费分时接口仅保留最近 5 个交易日)
- 未声明 / 走 tickflow → None (深历史)
preferences GET 带出该字段, 前端分时档位据此收窄 (浅源默认 5日, 深源默认 20日)。
"""
from __future__ import annotations
from types import SimpleNamespace
from app.api import settings
from app.services import preferences
def _mock_resolver(monkeypatch, provider, fallback, err=None):
monkeypatch.setattr(
"app.services.kline_sync._resolve_minute_provider",
lambda name: (provider, fallback, err),
)
def test_stocksdk_declares_five_day_history():
from app.plugins.stocksdk.provider import StockSDKProvider
assert StockSDKProvider.minute_history_days == 5
def test_history_days_from_custom_provider(monkeypatch):
"""自定义浅源 → 声明值; 前端据此只显示 1/5 日档。"""
_mock_resolver(monkeypatch, SimpleNamespace(minute_history_days=5), False)
monkeypatch.setattr(preferences, "get_minute_data_provider", lambda: "stocksdk")
assert settings._minute_history_days() == 5
def test_history_days_none_for_undeclared_provider(monkeypatch):
"""未声明的自定义源 → None (深历史基准)。"""
_mock_resolver(monkeypatch, SimpleNamespace(), False)
monkeypatch.setattr(preferences, "get_minute_data_provider", lambda: "my_source")
assert settings._minute_history_days() is None
def test_history_days_none_for_tickflow(monkeypatch):
"""tickflow (回退路径) → None (深历史)。"""
_mock_resolver(monkeypatch, None, True)
monkeypatch.setattr(preferences, "get_minute_data_provider", lambda: "tickflow")
assert settings._minute_history_days() is None
def test_history_days_none_when_resolver_fails(monkeypatch):
"""resolver 异常 (registry 损坏) → 降级 None, 不抛 500。"""
_mock_resolver(monkeypatch, None, True, err="registry broken")
monkeypatch.setattr(preferences, "get_minute_data_provider", lambda: "stocksdk")
assert settings._minute_history_days() is None
def test_preferences_get_includes_history_days(monkeypatch):
"""GET /preferences 响应包含 minute_history_days 字段。"""
_mock_resolver(monkeypatch, SimpleNamespace(minute_history_days=5), False)
payload = settings.get_preferences()
assert payload["minute_history_days"] == 5
assert "minute_data_provider" in payload
+101 -12
View File
@@ -34,7 +34,8 @@ class _FakeCapSet:
def has(self, cap) -> bool:
from app.tickflow.capabilities import Cap
return self._has and cap == Cap.INTRADAY_BATCH
# 服务门控挂 INTRADAY_UNIVERSE, 修复轮用 INTRADAY_BATCH — 两者同档, 一起授/不授
return self._has and cap in (Cap.INTRADAY_BATCH, Cap.INTRADAY_UNIVERSE)
class _FakeAppState:
@@ -160,13 +161,101 @@ def test_run_round_records_error_when_burst_empty(tmp_path, monkeypatch):
assert st["last_requests"] == 3
# ── 两段式模式选择: 冷启动全天 → 稳态增量 ───────────────────────────
def _patch_round(monkeypatch, *, lag, inc_df, burst_df):
calls: dict = {"modes": []}
monkeypatch.setattr(
minute_refresh.MinuteRefreshService, "_today_coverage_lag_minutes",
lambda self: lag, raising=True,
)
monkeypatch.setattr(
"app.services.kline_sync.fetch_intraday_universe_increment",
lambda *a, **k: (calls["modes"].append("increment"), (inc_df, 1))[1],
)
monkeypatch.setattr(
"app.services.kline_sync.fetch_intraday_full_market_burst",
lambda symbols, capset, *, count=300: (calls["modes"].append("full"), (burst_df, 28))[1],
)
monkeypatch.setattr(
"app.services.kline_sync._write_minute_partition",
lambda df, minute_dir: df.height,
)
return calls
def _inc_df():
return pl.DataFrame({
"symbol": ["600000.SH", "000001.SZ"],
"datetime": [datetime(2026, 8, 25, 10, 0)] * 2,
"open": [10.0] * 2, "high": [10.5] * 2, "low": [9.9] * 2, "close": [10.2] * 2,
"volume": [1000.0] * 2, "amount": [10200.0] * 2,
})
def _full_df():
return _inc_df()
def test_cold_start_no_local_data_uses_full_mode(tmp_path, monkeypatch):
"""当日无数据 (lag=None, 如 10 点冷启动) → 全天修复轮。"""
svc = _svc(tmp_path, monkeypatch)
calls = _patch_round(monkeypatch, lag=None, inc_df=_inc_df(), burst_df=_full_df())
svc._run_round()
assert calls["modes"] == ["full"]
st = svc.status()
assert st["last_mode"] == "full"
assert st["last_rows"] == 2
def test_healthy_coverage_uses_increment_mode(tmp_path, monkeypatch):
"""当日覆盖新鲜 (lag ≤ 3 分钟) → universe 单请求增量, 不打 burst。"""
svc = _svc(tmp_path, monkeypatch)
calls = _patch_round(monkeypatch, lag=0.2, inc_df=_inc_df(), burst_df=_full_df())
svc._run_round()
assert calls["modes"] == ["increment"]
st = svc.status()
assert st["last_mode"] == "increment"
assert st["last_requests"] == 1
assert st["last_symbols"] == 2
assert st["last_rows"] == 2
def test_stale_coverage_beyond_bar_headroom_falls_back_to_full(tmp_path, monkeypatch):
"""覆盖滞后超过 3 分钟 (超过 universe 3 根余量) → 全天修复轮。"""
svc = _svc(tmp_path, monkeypatch)
calls = _patch_round(monkeypatch, lag=5.0, inc_df=_inc_df(), burst_df=_full_df())
svc._run_round()
assert calls["modes"] == ["full"]
def test_consecutive_empty_rounds_escalate_to_full(tmp_path, monkeypatch):
"""universe 连续 2 轮空返回 → 第 3 轮自动升级全天修复 (自愈)。"""
svc = _svc(tmp_path, monkeypatch)
calls = _patch_round(
monkeypatch,
lag=0.2,
inc_df=pl.DataFrame(), # 增量恒空 (模拟 universe 端点持续异常)
burst_df=_full_df(),
)
svc._run_round()
svc._run_round()
assert calls["modes"] == ["increment", "increment"]
assert svc.status()["rounds"] == 0
svc._run_round()
assert calls["modes"] == ["increment", "increment", "full"]
assert svc.status()["last_mode"] == "full"
def test_status_reports_gate_reason_when_stopped(tmp_path, monkeypatch):
svc = _svc(tmp_path, monkeypatch, enabled=False)
st = svc.status()
assert st["enabled"] is False
assert st["running"] is False
assert st["gate_reason"] == "disabled"
assert st["interval_seconds"] == 60
assert st["interval_seconds"] == 6
# ── 偏好 ────────────────────────────────────────────────────────────
@@ -175,27 +264,27 @@ def test_status_reports_gate_reason_when_stopped(tmp_path, monkeypatch):
def test_refresh_preferences_defaults_and_clamp(tmp_path, monkeypatch):
_isolated_prefs(tmp_path, monkeypatch)
assert preferences.get_minute_refresh_enabled() is False
assert preferences.get_minute_refresh_interval() == 60
preferences.save({"minute_refresh_interval": 5})
assert preferences.get_minute_refresh_interval() == 60 # 下限
assert preferences.get_minute_refresh_interval() == 6
preferences.save({"minute_refresh_interval": 1})
assert preferences.get_minute_refresh_interval() == 3 # 下限
preferences.save({"minute_refresh_interval": 999})
assert preferences.get_minute_refresh_interval() == 300 # 上限
preferences.save({"minute_refresh_interval": 90})
assert preferences.get_minute_refresh_interval() == 90
preferences.save({"minute_refresh_interval": 15})
assert preferences.get_minute_refresh_interval() == 15
def test_realtime_monitor_config_owns_refresh_keys(tmp_path, monkeypatch):
"""盘中增量配置归属实时监控端点 (set_realtime_monitor_config), 并 clamp 到 [60,300]。"""
"""盘中增量配置归属实时监控端点 (set_realtime_monitor_config), 并 clamp 到 [3,300]。"""
_isolated_prefs(tmp_path, monkeypatch)
saved = preferences.set_realtime_monitor_config({
"minute_refresh_enabled": True,
"minute_refresh_interval": 10, # 越界 → clamp 到下限
"minute_refresh_interval": 1, # 越界 → clamp 到下限
})
assert saved["minute_refresh_enabled"] is True
assert saved["minute_refresh_interval"] == 60
assert saved["minute_refresh_interval"] == 3
saved = preferences.set_realtime_monitor_config({"minute_refresh_interval": 400})
assert saved["minute_refresh_interval"] == 300
saved = preferences.set_realtime_monitor_config({"minute_refresh_interval": 120})
assert saved["minute_refresh_interval"] == 120
saved = preferences.set_realtime_monitor_config({"minute_refresh_interval": 6})
assert saved["minute_refresh_interval"] == 6
def test_status_endpoint_without_service():
@@ -0,0 +1,185 @@
"""分钟 K datetime 北京墙钟契约测试。
契约 (CONTRIBUTING §3.3): kline_minute.datetime 必须是北京墙钟 naive。
守卫 _enforce_minute_beijing_wallclock 在两个源头入口强制:
- _normalize_minute (TickFlow 帧, timestamp 毫秒为 UTC 基准)
- _try_custom_minute (插件/自定义源帧)
覆盖: 显式转换 / 北京墙钟直通 / UTC 特征自愈 +8 / tz-aware 换算 /
fail-closed 拒收 / 路由级契约违规回退 TickFlow。
"""
from __future__ import annotations
from datetime import UTC, datetime
from unittest.mock import MagicMock
import polars as pl
import pytest
from app.services import kline_sync
def _minute_frame(datetimes: list, symbol: str = "600519.SH") -> pl.DataFrame:
n = len(datetimes)
return pl.DataFrame({
"symbol": [symbol] * n,
"datetime": datetimes,
"open": [10.0] * n,
"high": [10.5] * n,
"low": [9.5] * n,
"close": [10.2] * n,
"volume": [100.0] * n,
"amount": [1020.0] * n,
})
def _beijing_day() -> list[datetime]:
"""一个正常交易日墙钟样本: 开盘/午盘首/收盘。"""
return [
datetime(2026, 1, 15, 9, 30),
datetime(2026, 1, 15, 13, 0),
datetime(2026, 1, 15, 15, 0),
]
# ---------- TickFlow 路径: timestamp 毫秒 (UTC 基准) → 北京墙钟 ----------
def test_tickflow_timestamp_normalizes_to_beijing_wallclock():
"""09:30 北京 = 01:30 UTC; SDK 帧 timestamp 毫秒归一后必须回到 09:30。"""
ts_ms = [
int(datetime(2026, 1, 15, 1, 30, tzinfo=UTC).timestamp() * 1000), # 09:30 北京
int(datetime(2026, 1, 15, 5, 0, tzinfo=UTC).timestamp() * 1000), # 13:00 北京
int(datetime(2026, 1, 15, 7, 0, tzinfo=UTC).timestamp() * 1000), # 15:00 北京
]
df = pl.DataFrame({
"symbol": ["600519.SH"] * 3,
"timestamp": ts_ms,
"open": [10.0] * 3, "high": [10.5] * 3,
"low": [9.5] * 3, "close": [10.2] * 3,
"volume": [100.0] * 3, "amount": [1020.0] * 3,
})
out = kline_sync._normalize_minute(df)
assert out["datetime"].to_list() == _beijing_day()
def test_tickflow_timestamp_partial_day_lunch_unaffected():
"""午间 11:30 (03:30 UTC) 同样正确归一, 不被误判为越界。"""
df = pl.DataFrame({
"symbol": ["000001.SZ"],
"timestamp": [int(datetime(2026, 1, 15, 3, 30, tzinfo=UTC).timestamp() * 1000)],
"open": [1.0], "high": [1.0], "low": [1.0], "close": [1.0],
"volume": [1.0], "amount": [1.0],
})
out = kline_sync._normalize_minute(df)
assert out["datetime"].to_list() == [datetime(2026, 1, 15, 11, 30)]
# ---------- 守卫: 各口径分类 ----------
def test_guard_beijing_naive_passthrough():
df = _minute_frame(_beijing_day())
out = kline_sync._enforce_minute_beijing_wallclock(df, source="t")
assert out["datetime"].to_list() == _beijing_day()
def test_guard_utc_naive_selfhealed_plus8():
"""01:30/05:00/07:00 (UTC 墙钟特征) → 自动 +8 → 09:30/13:00/15:00。"""
df = _minute_frame([
datetime(2026, 1, 15, 1, 30),
datetime(2026, 1, 15, 5, 0),
datetime(2026, 1, 15, 7, 0),
])
out = kline_sync._enforce_minute_beijing_wallclock(df, source="t")
assert out["datetime"].to_list() == _beijing_day()
def test_guard_selfheal_is_idempotent():
df = _minute_frame([datetime(2026, 1, 15, 1, 30), datetime(2026, 1, 15, 3, 0)])
once = kline_sync._enforce_minute_beijing_wallclock(df, source="t")
twice = kline_sync._enforce_minute_beijing_wallclock(once, source="t")
assert once["datetime"].to_list() == twice["datetime"].to_list()
def test_guard_tzaware_utc_converted():
"""tz-aware UTC 01:30 → 北京墙钟 09:30 (naive)。"""
df = _minute_frame([datetime(2026, 1, 15, 1, 30, tzinfo=UTC)])
out = kline_sync._enforce_minute_beijing_wallclock(df, source="t")
assert out["datetime"].dtype == pl.Datetime("us")
assert out["datetime"].to_list() == [datetime(2026, 1, 15, 9, 30)]
def test_guard_tzaware_shanghai_converted():
"""tz-aware +08:00 09:30 → 北京墙钟 09:30 (naive), 数值不变。"""
df = _minute_frame([datetime(2026, 1, 15, 9, 30, tzinfo=_shanghai_tz())])
out = kline_sync._enforce_minute_beijing_wallclock(df, source="t")
assert out["datetime"].to_list() == [datetime(2026, 1, 15, 9, 30)]
def _shanghai_tz():
from zoneinfo import ZoneInfo
return ZoneInfo("Asia/Shanghai")
def test_guard_unrecognized_convention_fails_closed():
"""21:30/22:15 (境外墙钟特征) 既非北京时段也非 UTC 平移 → 拒收。"""
df = _minute_frame([datetime(2026, 1, 15, 21, 30), datetime(2026, 1, 15, 22, 15)])
with pytest.raises(ValueError, match="口径无法识别"):
kline_sync._enforce_minute_beijing_wallclock(df, source="t")
def test_guard_all_null_datetimes_passthrough():
"""全 null datetime 维持原行为 (下游落盘过滤), 不误伤。"""
df = _minute_frame([None, None])
out = kline_sync._enforce_minute_beijing_wallclock(df, source="t")
assert out.height == 2
assert out["datetime"].null_count() == 2
def test_guard_string_datetimes_classified_after_parse():
"""trade_time 字符串路径: 先解析再分类 (UTC 特征串同样自愈)。"""
df = pl.DataFrame({
"symbol": ["600519.SH"],
"trade_time": ["2026-01-15 01:30:00"],
"open": [1.0], "high": [1.0], "low": [1.0], "close": [1.0],
"volume": [1.0], "amount": [1.0],
}).rename({"trade_time": "datetime"})
out = kline_sync._enforce_minute_beijing_wallclock(df, source="t")
assert out["datetime"].to_list() == [datetime(2026, 1, 15, 9, 30)]
# ---------- 路由级: 自定义源契约违规 → 回退 TickFlow ----------
def _setup_custom_provider(monkeypatch, provider: object) -> None:
monkeypatch.setattr(kline_sync.preferences, "get_minute_data_provider", lambda: "mock_src")
monkeypatch.setattr("app.data_providers.custom.provider_has_dataset", lambda name, ds: True)
monkeypatch.setattr("app.data_providers.custom.get_provider", lambda name: provider)
def test_custom_provider_utc_frame_selfhealed(monkeypatch):
"""插件返回 UTC 墙钟帧 → 路由层守卫 +8 后下发, 不回退。"""
mock_provider = MagicMock()
mock_provider.get_minute = MagicMock(return_value=_minute_frame(
[datetime(2026, 1, 15, 1, 30), datetime(2026, 1, 15, 5, 0)]))
_setup_custom_provider(monkeypatch, mock_provider)
df, fallback = kline_sync._try_custom_minute(
["600519.SH"], datetime(2026, 1, 15, 9, 25), datetime(2026, 1, 15, 15, 5),
asset_type="stock",
)
assert fallback is False
assert df["datetime"].to_list() == [datetime(2026, 1, 15, 9, 30), datetime(2026, 1, 15, 13, 0)]
def test_custom_provider_garbage_datetime_falls_back(monkeypatch):
"""插件返回无法识别口径 → fail-closed 回退 TickFlow。"""
mock_provider = MagicMock()
mock_provider.get_minute = MagicMock(return_value=_minute_frame(
[datetime(2026, 1, 15, 21, 30)]))
_setup_custom_provider(monkeypatch, mock_provider)
df, fallback = kline_sync._try_custom_minute(
["600519.SH"], datetime(2026, 1, 15, 9, 25), datetime(2026, 1, 15, 15, 5),
asset_type="stock",
)
assert fallback is True
assert df is None
+3 -3
View File
@@ -2493,15 +2493,15 @@ wheels = [
[[package]]
name = "tickflow"
version = "0.1.24"
version = "0.1.25"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "httpx" },
{ name = "typing-extensions" },
]
sdist = { url = "https://files.pythonhosted.org/packages/aa/56/d911f7d03363a06f69838878df4f92dd01235c899431197f94fb4c0e36ad/tickflow-0.1.24.tar.gz", hash = "sha256:13f6464a9dd1bdf98a312bc8c313fcc24ad0ad83c16d4f83b76bee83df6e11df", size = 36867, upload-time = "2026-06-20T02:39:35.891Z" }
sdist = { url = "https://files.pythonhosted.org/packages/d1/c9/facd0cd7568ea3c7dcfb0d52ded7d0b52ad871ca45fa97d7df0925b26091/tickflow-0.1.25.tar.gz", hash = "sha256:a86929bc99167014567c3d8b99af1c2f498922ad476bb4507edb57164c576d04", size = 37357, upload-time = "2026-08-29T04:43:27.255Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/4c/9d/6c03706054f3bcca8a7113a60258e1a52a762077a555d52a5a84e4f6895f/tickflow-0.1.24-py3-none-any.whl", hash = "sha256:e898867b0e3e668618135c78e3a367542f81b7a289567335d298c707452e5f42", size = 43031, upload-time = "2026-06-20T02:39:34.307Z" },
{ url = "https://files.pythonhosted.org/packages/ea/37/f80b8f6e435f1825ea8384b1605a01c323998f00b1c112c59e72a9d7dd9a/tickflow-0.1.25-py3-none-any.whl", hash = "sha256:69687235e85b44eae077262325ddef15e82d5adafda075945612f8cf0815a44c", size = 43510, upload-time = "2026-08-29T04:43:25.785Z" },
]
[package.optional-dependencies]
+6 -4
View File
@@ -192,7 +192,8 @@ export interface AiStockReport {
// ===== Kline =====
export interface MinuteKlineRow {
datetime: string
open: number
/** 分钟开盘价; 部分数据源(stock-sdk 历史日)无真实分钟 open, 为 null */
open: number | null
high: number
low: number
close: number
@@ -1407,7 +1408,7 @@ export interface CapabilityRoute {
id: string
label: string
desc: string
field: ProviderField
field: ProviderField | null // null = 不可路由能力 (仅 TickFlow 提供)
default: string
tf_tier: string // TickFlow 所需最低订阅档位
tf_available: boolean // 当前 TickFlow 档位是否提供该能力
@@ -1508,12 +1509,13 @@ export interface Preferences {
daily_data_provider?: string
adj_factor_provider?: string
minute_data_provider?: string
/** 分钟源 1 分钟历史深度(交易日); null/缺省 = 深历史 (如 tickflow)。分时档位据此收窄 */
minute_history_days?: number | null
depth5_data_provider?: string
realtime_data_provider?: string
financial_data_provider?: string
data_source_job_timeout_s: number
data_source_long_job_timeout_s: number
realtime_watchlist_symbols?: string[]
realtime_pull_stock?: boolean
realtime_pull_etf?: boolean
realtime_pull_index?: boolean
@@ -1704,7 +1706,7 @@ export const api = {
}),
}),
/** 盘中分钟增量刷新服务状态 (Expert 专有) */
/** 全量分钟 (盘中全市场分钟落盘) 服务状态 (TickFlow Expert 专有) */
minuteRefreshStatus: () =>
request<{
available: boolean
+1
View File
@@ -10,6 +10,7 @@ export const CAP_LABELS: Record<string, { name: string; hint: string }> = {
'kline.daily.batch': { name: '日 K(批量)', hint: '一次拿多只股票的日 K — 选股 / 信号扫描 必需' },
'kline.minute.by_symbol': { name: '分钟 K(按标的)', hint: '单股 1m/5m/15m/30m/60m K 线' },
'kline.minute.batch': { name: '分钟 K(批量)', hint: '多股分钟 K' },
'intraday.universe': { name: '全量分钟', hint: '标的池单请求拉全市场当日分钟K (盘中增量落盘, Expert 专有)' },
'depth5': { name: '五档盘口', hint: '买卖五档报价' },
'depth5.batch': { name: '五档盘口(批量)', hint: '批量买卖五档快照' },
+17 -68
View File
@@ -1,5 +1,4 @@
import { useState, useCallback, useEffect, createContext, useContext } from 'react'
import { Link } from 'react-router-dom'
import { useQueryClient, useMutation, useQuery } from '@tanstack/react-query'
import {
Activity,
@@ -51,16 +50,14 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
const { data: intervalData } = useQuoteInterval()
const updateInterval = useUpdateQuoteInterval()
const toggleQuote = useToggleRealtimeQuotes()
// 实时模式以 quote_status 为准 (数据源无关): watchlist=自选实时 / full_market=全市场 / none=不可用
const quoteMode = quoteStatus?.mode ?? 'none'
const isWatchlistMode = quoteMode === 'watchlist'
// 实时模式以 quote_status 为准 (数据源无关): full_market=全市场 / none=不可用
const realtimeEnabled = prefs?.realtime_quotes_enabled ?? false
// 分时图实时刷新间隔 (秒), 与后端 [3,60] clamp 对齐; 默认 6
const intradayInterval = prefs?.minute_intraday_refresh_interval ?? 6
// 滑块本地草稿: 拖动时即时反馈, 停顿 2s 后落库 (与行情轮询滑块一致)
const [intradayIntervalDraft, setIntradayIntervalDraft] = useState(intradayInterval)
// 盘中分钟增量 (Expert 专有): 间隔 (秒), 与后端 [60,300] clamp 对齐; 默认 60
const minuteRefreshInterval = prefs?.minute_refresh_interval ?? 60
// 盘中分钟增量 (Expert 专有): 间隔 (秒), 与后端 [3,300] clamp 对齐; 默认 6
const minuteRefreshInterval = prefs?.minute_refresh_interval ?? 6
const [minuteRefreshIntervalDraft, setMinuteRefreshIntervalDraft] = useState(minuteRefreshInterval)
// 盘中增量服务状态 (15s 轮询; 无服务时 available=false)
const refreshStatus = useQuery({
@@ -71,8 +68,9 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
const refreshPages = prefs?.sse_refresh_pages ?? {}
const limitLadderMonitor = prefs?.limit_ladder_monitor_enabled ?? false
const hasDepth = !!caps?.capabilities?.['depth5.batch']
// 盘中分钟增量 = intraday.batch 独立能力 (Expert 专有), 与盘后同步的 minute.batch 分属不同限流池
const hasIntradayBatchCap = !!caps?.capabilities?.['intraday.batch']
// 全量分钟 = intraday.universe 能力 (TickFlow Expert 专有): 标的池单请求拉全市场当日分钟,
// 修复轮的 intraday.batch 与其同档, 见后端 minute_refresh 服务
const hasFullMinuteCap = !!caps?.capabilities?.['intraday.universe']
const rs = refreshStatus.data
// 新建监控规则时默认勾选的推送渠道 (全局默认值数组, 单条规则可独立修改)
const webhookDefaultChannels = prefs?.webhook_default_channels ?? []
@@ -120,15 +118,6 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
setBotIdDraft(wecomBotId)
setBotSecretDraft(wecomBotSecret)
}, [wecomBotId, wecomBotSecret])
const watchlistSymbols = prefs?.realtime_watchlist_symbols ?? []
const watchlist = useQuery({
queryKey: QK.watchlist,
queryFn: () => api.watchlistList(),
enabled: isWatchlistMode && watchlistSymbols.length > 0,
})
const watchlistNameBySymbol = new Map(
(watchlist.data?.symbols ?? []).map(row => [row.symbol, row.name] as const),
)
const save = useCallback(async (cfg: Record<string, unknown>) => {
try {
@@ -318,7 +307,7 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
<div className="min-w-0">
<div className="text-sm text-foreground"></div>
<div className="text-[11px] text-muted">
{isWatchlistMode ? '每轮拉取自选股实时行情的时间间隔' : '每轮拉取全市场行情的时间间隔'}
</div>
</div>
<span className="text-[11px] font-mono text-foreground shrink-0 tabular-nums">
@@ -342,43 +331,6 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
</div>
</Card>
{isWatchlistMode && (
<Card icon={Activity} title="自选股实时">
<div className="mb-3 rounded-btn border border-accent/25 bg-accent/10 px-3 py-2 text-xs font-medium leading-snug text-accent">
5 6
</div>
{watchlistSymbols.length > 0 ? (
<div className="space-y-1.5">
{watchlistSymbols.map(symbol => {
const name = watchlistNameBySymbol.get(symbol)
return (
<div key={symbol} className="flex items-center justify-between rounded-btn bg-base/50 border border-border px-2 py-1.5">
<div className="min-w-0 flex items-baseline gap-1.5">
<span className="text-xs font-mono text-foreground">{symbol}</span>
{name && <span className="truncate text-[11px] text-secondary">{name}</span>}
</div>
<span className="text-[10px] text-muted shrink-0"></span>
</div>
)
})}
</div>
) : (
<div className="rounded-btn border border-border bg-base/40 px-3 py-3 text-xs text-muted">
</div>
)}
<div className="mt-2 flex items-center justify-between gap-3">
<span className="text-[10px] text-muted"> {watchlistSymbols.length}/5 </span>
<Link
to="/watchlist"
className="px-3 py-1 rounded-btn bg-elevated text-secondary text-xs font-medium hover:text-foreground transition-colors"
>
</Link>
</div>
</Card>
)}
{!isWatchlistMode && (
<Card icon={Wifi} title="页面实时刷新">
<p className="text-xs text-secondary mb-4">
SSE
@@ -396,7 +348,6 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
))}
</div>
</Card>
)}
{/* 自选列表分时图实时刷新 (默认关闭, 开启后盘中按设定间隔轮询刷新分时数据) */}
<Card icon={Activity} title="分时图刷新" anchor="intraday-refresh">
@@ -435,7 +386,6 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
</div>
</Card>
{!isWatchlistMode && (
<Card icon={BarChart3} title="左侧菜单指数">
<p className="text-xs text-secondary mb-4">
@@ -460,7 +410,6 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
/>
</div>
</Card>
)}
</div>
{/* ========== 右列 ========== */}
@@ -508,26 +457,26 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
)}
</Card>
{/* 盘中分钟增量落盘 (Expert 专有): 交易时段常驻服务, intraday.batch 独立配额 */}
<Card icon={Zap} title="盘中分钟增量" anchor="minute-refresh">
{/* 全量分钟 (TickFlow Expert 专有): 盘中全市场分钟落盘, intraday.universe 单请求增量 */}
<Card icon={Zap} title="全量分钟" anchor="minute-refresh">
<ToggleRow
label="盘中分钟增量落盘"
label="全量分钟落盘"
desc={
!hasIntradayBatchCap ? '需要日内分时批量能力 (Expert)'
!hasFullMinuteCap ? '需要全量分钟能力 (TickFlow Expert)'
: rs?.custom_provider_active ? '已配置自定义分钟源, 盘中增量由插件自管'
: rs?.running ? (rs?.in_trading_hours ? '服务运行中' : '运行中 · 非连续竞价时段暂停')
: '已关闭'
}
checked={prefs?.minute_refresh_enabled ?? false}
onChange={(v) => save({ minute_refresh_enabled: v })}
disabled={!hasIntradayBatchCap || !!rs?.custom_provider_active}
disabled={!hasFullMinuteCap || !!rs?.custom_provider_active}
/>
<div className="mt-3 pt-3 border-t border-border">
<div className="flex items-center justify-between gap-4 py-1">
<div className="min-w-0">
<div className="text-sm text-foreground"></div>
<div className="text-[11px] text-muted">
; 60s intraday.batch
K增量落盘的间隔; , /
</div>
</div>
<span className="text-[11px] font-mono text-foreground shrink-0 tabular-nums">
@@ -537,16 +486,16 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
<div className="flex items-center gap-3 mt-2">
<input
type="range"
min={60}
min={3}
max={300}
step={30}
step={3}
value={minuteRefreshIntervalDraft}
disabled={!hasIntradayBatchCap}
disabled={!hasFullMinuteCap}
onChange={(e) => setMinuteRefreshIntervalDraft(parseInt(e.target.value, 10))}
className="flex-1 h-1 accent-accent cursor-pointer disabled:opacity-40 disabled:cursor-not-allowed"
/>
<span className="text-[10px] text-muted shrink-0">
{minuteRefreshIntervalDraft !== minuteRefreshInterval ? '2秒后保存' : '60s — 300s'}
{minuteRefreshIntervalDraft !== minuteRefreshInterval ? '2秒后保存' : '3s — 300s'}
</span>
</div>
{rs?.available && rs.rounds != null && rs.rounds > 0 && (
+1
View File
@@ -66,6 +66,7 @@ expert:
kline.minute.by_symbol: { rpm: 120, batch: 1 }
intraday: { rpm: 120, batch: 1 }
intraday.batch: { rpm: 60, batch: 200 }
intraday.universe: { rpm: 20 } # 全量分钟增量: 标的池单请求拉全市场; rpm=20 对应服务 3s 下限
depth5: { rpm: 120, batch: 1 } # 按标的查(单只):官方 120rpm/1
depth5.batch: { rpm: 60, batch: 200 } # 批量查(新增):官方 60rpm/200
adj_factor: { rpm: 120, batch: 200 }