Compare commits

..
Author SHA1 Message Date
a836cda1b2 feat(HEL-421): 回补历史日历和指数并标记区间不完整
Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
2026-09-02 22:16:05 +08:00
总工 25ff6bbe06 fix(HEL-417): 配置根日志让影子对比报告落入容器日志 2026-09-02 21:53:19 +08:00
5085cacf0d fix(HEL-412): 刷新降级不再整次失败,并补齐准备中提示
手动刷新与自动补跑共用可用数据判定:日线推算或上一交易日快照记为部分/准备中成功,避免前端误报刷新失败。HTTP JSON 解析错误不再把请求正文写入日志。

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
2026-09-02 18:08:40 +08:00
38 changed files with 1312 additions and 93 deletions
+3 -1
View File
@@ -54,7 +54,9 @@ background scheduler
feature repository mixins; do not add feature queries to it.
- `backend/jobs/` owns job definitions, locks, retries, idempotency, and persisted run state.
`backend/jobs/service.py` is the application-facing owner of scheduler start/stop, manual
refresh submission, and periodic refresh coordination.
refresh submission, and periodic refresh coordination. `backend/jobs/refresh.py` owns
whether a dashboard payload is a usable refresh result versus a failed job, and whether
after-hours official catch-up is due.
- `backend/llm/` owns model selection, membership/quota checks, fallback, provider transport,
streaming rules, and call audit. Feature agents only prepare messages and interpret
feature-specific results.
+2 -2
View File
@@ -16,7 +16,7 @@
- **题材库 / 人气热榜 / 龙虎榜**:题材成分、双榜人气、席位与游资档案
- **智能选股**(会员):六阶段策略、精选策略库、自然语言编译为受控公式后的确定性筛选与滚动回测;候选需手动加入后才进入五交易日跟踪
- **问师**(会员):按选定的游资思维 Skill 单师对话;新增公开角色时在 `游资skills` 下增加含 `SKILL.md` 的目录,并在 `游资skills/mentor_catalog.json` 登记。管理员私有角色放在 `data/private-mentor-skills`(不进 Git / 镜像)
- **问天**(会员,冻结区,勿改代码):观势 / 观气 / 观心。卦象、干支、节气与气机由本地程序确定性计算,大模型只负责文字解释
- **问天**(会员):观势 / 观气 / 观心。卦象、干支、节气与气机由本地程序确定性计算,大模型只负责文字解释。此前仅冻结过界面视觉方案,现已解冻;问天可纳入后续数据与功能迁移,本阶段不主动重做视觉。
- **我的复盘**:手工交易日志、每日复盘、提醒中心与复盘助手;不接券商、不自动下单
全局能力:日间 / 夜间主题、股票代码悬停预览日 K 与分时、`Ctrl + K` 全局搜索。图表数据不写入主行情,也不参与情绪、选股或问天计算。
@@ -132,7 +132,7 @@ compose.yaml
- 本项目是个人研究与复盘工具,全部数据、指标、候选与文字分析均不构成投资建议、证券推荐或买卖要约。
- 不接券商、不代为下单。交易日志只做手工记录与统计,不代表实际成交。
- 情绪温度、阶段判定、连板梯队、策略筛选等均为基于公开数据的统计与规则计算,不预测走势,不保证收益。
- 「问天」属于传统文化视角的观察工具,不具备预测功能,不得作为投资依据;该模块为冻结区,不要改其代码
- 「问天」属于传统文化视角的观察工具,不具备预测功能,不得作为投资依据。问天不是永久冻结区:此前只冻结过界面视觉方案,现已解冻,后续数据与功能迁移可以纳入
- 行情来自第三方接口,可能延迟、缺失或口径调整;不可用时页面会明确提示,请以交易所与券商正式披露为准。
- 不要把服务端口直接暴露到公网。不要把 Token、密码、密钥、数据库或 `.env` 提交进 Git。
- 股市有风险,入市需谨慎。投资决策及其后果由使用者本人承担。
+12
View File
@@ -1,11 +1,23 @@
from __future__ import annotations
import argparse
import logging
from http.server import ThreadingHTTPServer
from typing import Any
def configure_logging() -> None:
"""让 INFO 级结构化日志(含 datahub 影子对比报告)落到容器日志。"""
if logging.getLogger().handlers:
return
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(name)s %(message)s",
)
def main(handler_class: type[Any] | None = None, service: Any | None = None) -> None:
configure_logging()
if handler_class is None or service is None:
from backend.application import RequestHandler, SERVICE
+6
View File
@@ -25,6 +25,7 @@ EMPTY_FAIL_DATASETS = {"stocks", "daily", "index_daily", "valuation", "moneyflow
def looks_like_heaven(module_name: str, filename: str = "") -> bool:
"""问天调用栈识别。问天未永久冻结,只是本阶段仍走旧 Tushare 链路。"""
path = filename.replace("\\", "/")
return module_name.startswith("backend.features.heaven") or "/features/heaven/" in path
@@ -95,6 +96,7 @@ class DatahubBridge:
legacy_query: Callable[..., list[dict[str, Any]]],
) -> list[dict[str, Any]]:
dataset = API_TO_DATASET.get(api_name)
# 问天允许后续纳入 datahub;首批只读接入仍保持旧链路,避免误切。
if not dataset or self.heaven_guard():
return legacy_query(api_name, params, fields)
flags = self.settings.flags(dataset)
@@ -206,6 +208,10 @@ class DatahubBridge:
raise DatahubError("STALE", f"{dataset} data is stale")
if dataset in EMPTY_FAIL_DATASETS and not rows:
raise DatahubError("EMPTY", f"{dataset} returned no rows")
coverage = meta.get("coverage") if isinstance(meta.get("coverage"), dict) else {}
if meta.get("incomplete") is True or coverage.get("complete") is False:
missing = coverage.get("missing_count")
raise DatahubError("INCOMPLETE", f"{dataset} range is incomplete missing={missing}")
def _require_fresh(self, response: DatahubResponse, dataset: str) -> DatahubResponse:
self._validate_usable(dataset, list(response.data or []) if isinstance(response.data, list) else [], response)
+79 -10
View File
@@ -79,6 +79,8 @@ class MarketServiceMixin:
if not force:
snapshot = self.database.get_snapshot(normalized_date)
if snapshot and str((snapshot.get("meta") or {}).get("source") or "") != "demo":
if self._should_retry_incomplete_snapshot(snapshot, normalized_date):
return self.sync_dashboard(normalized_date)
snapshot = copy.deepcopy(snapshot)
if normalized_date != now.strftime("%Y%m%d"):
snapshot.setdefault("meta", {}).update(
@@ -97,6 +99,8 @@ class MarketServiceMixin:
"dashboard_request_v1", normalized_date
)
if resolved and str((resolved.get("meta") or {}).get("source") or "") != "demo":
if self._should_retry_incomplete_snapshot(resolved, normalized_date):
return self.sync_dashboard(normalized_date)
resolved = copy.deepcopy(resolved)
resolved.setdefault("meta", {})["requested_date"] = self._display_compact_date(
normalized_date
@@ -138,6 +142,68 @@ class MarketServiceMixin:
def _display_compact_date(compact: str) -> str:
return f"{compact[:4]}-{compact[4:6]}-{compact[6:8]}"
@staticmethod
def _chinese_month_day(value: str) -> str:
compact = str(value or "").replace("-", "").replace("/", "")
if len(compact) < 8 or not compact[:8].isdigit():
return "最近可用交易日"
return f"{int(compact[4:6])}{int(compact[6:8])}"
@classmethod
def _preparing_display_notice(cls, actual_date: str, requested_date: str) -> str:
shown = cls._chinese_month_day(actual_date)
requested = str(requested_date or "").replace("-", "")
if requested == date.today().strftime("%Y%m%d"):
return f"今日数据正在准备,当前展示 {shown}"
return f"所选日期数据尚未到齐,当前展示 {shown}"
@staticmethod
def _snapshot_age_seconds(meta: dict[str, Any]) -> float:
raw = str(meta.get("updated_at") or "")
if not raw:
return 10**9
try:
updated_at = datetime.fromisoformat(raw)
except ValueError:
return 10**9
now = datetime.now().astimezone()
if updated_at.tzinfo is None:
updated_at = updated_at.replace(tzinfo=now.tzinfo)
return (now - updated_at.astimezone(now.tzinfo)).total_seconds()
def _should_retry_incomplete_snapshot(
self, snapshot: dict[str, Any], requested_date: str
) -> bool:
if requested_date != date.today().strftime("%Y%m%d"):
return False
meta = snapshot.get("meta") or {}
incomplete = (
meta.get("limit_data_source") == "derived"
or bool(meta.get("carried_forward"))
or str(meta.get("trade_date") or "").replace("-", "") != requested_date
)
return incomplete and self._snapshot_age_seconds(meta) >= 60
def _annotate_data_status(self, dashboard: dict[str, Any]) -> dict[str, Any]:
meta = dashboard.setdefault("meta", {})
notice = str(meta.get("notice") or "")
requested = str(meta.get("requested_date") or "").replace("-", "")
actual = str(meta.get("trade_date") or "").replace("-", "")
if meta.get("limit_data_source") == "derived" and not meta.get("carried_forward"):
meta["data_status"] = "partial"
meta["display_notice"] = notice or "部分正式数据尚未到齐,当前展示日线推算结果"
elif meta.get("carried_forward"):
if "非交易日" in notice or "盘前" in notice:
meta["data_status"] = "carried"
meta["display_notice"] = notice
else:
meta["data_status"] = "preparing"
meta["display_notice"] = self._preparing_display_notice(actual, requested)
else:
meta["data_status"] = "official"
meta.setdefault("display_notice", "")
return dashboard
def _carry_dashboard(
self, snapshot: dict[str, Any], requested_date: str, reason: str
) -> dict[str, Any]:
@@ -152,7 +218,7 @@ class MarketServiceMixin:
"notice": reason,
}
)
return carried
return self._annotate_data_status(carried)
def _realtime_snapshot_due(
self,
@@ -197,14 +263,14 @@ class MarketServiceMixin:
if not self.configured:
raise TushareError("公共行情尚未配置")
dashboard = self._tushare_client().dashboard(normalized_date)
if (dashboard.get("meta") or {}).get("limit_data_source") == "derived":
raise TushareError(
str((dashboard.get("meta") or {}).get("notice") or "官方涨跌停数据尚未返回")
meta = dashboard.setdefault("meta", {})
meta["source"] = source
meta["requested_date"] = self._display_compact_date(normalized_date)
if meta.get("limit_data_source") == "derived":
meta.setdefault(
"notice",
"涨跌停高级接口当日数据尚未更新,已使用日线数据推算。",
)
dashboard["meta"]["source"] = source
dashboard["meta"]["requested_date"] = self._display_compact_date(normalized_date)
dashboard = self._enrich_dashboard_sentiment(dashboard, normalized_date)
record_count = self._record_count(dashboard)
actual_date = normalize_date(
@@ -233,8 +299,11 @@ class MarketServiceMixin:
except TushareError as exc:
fallback = self.database.get_latest_real_snapshot(normalized_date)
if fallback:
actual = str((fallback.get("meta") or {}).get("trade_date") or "")
carried = self._carry_dashboard(
fallback, normalized_date, f"最新行情暂不可用,沿用最近收盘快照:{exc}"
fallback,
normalized_date,
self._preparing_display_notice(actual, normalized_date),
)
self.database.finish_sync(
sync_id, "fallback", self._record_count(carried), str(exc), "tushare"
@@ -1160,7 +1229,7 @@ class MarketServiceMixin:
"storage": "sqlite",
"cached": cached,
}
return result
return self._annotate_data_status(result)
@staticmethod
def _record_count(dashboard: dict[str, Any]) -> int:
+8 -1
View File
@@ -109,7 +109,14 @@ class HttpTransportMixin:
return {}
if length <= 0 or length > 65536:
raise ValueError("请求内容为空或过大。")
return json.loads(self.rfile.read(length).decode("utf-8"))
raw = self.rfile.read(length)
try:
payload = json.loads(raw.decode("utf-8"))
except (UnicodeDecodeError, json.JSONDecodeError):
raise ValueError("请求不是合法 JSON。") from None
if not isinstance(payload, dict):
raise ValueError("请求不是合法 JSON。")
return payload
def serve_static(self, request_path: str) -> None:
relative = unquote(request_path).lstrip("/") or "index.html"
+46
View File
@@ -0,0 +1,46 @@
from __future__ import annotations
from datetime import datetime, time as dt_time
def dashboard_has_usable_data(dashboard: dict[str, object]) -> bool:
if not isinstance(dashboard, dict) or dashboard.get("status") == "failed":
return False
meta = dashboard.get("meta") or {}
overview = dashboard.get("overview") or {}
if isinstance(meta, dict) and (meta.get("trade_date") or meta.get("carried_forward")):
return True
return bool(isinstance(overview, dict) and overview)
def verified_dashboard_result(dashboard: dict[str, object]) -> dict[str, object]:
"""Manual refresh and automatic catch-up share this rule.
Derived limit lists or a previous usable snapshot are not whole-job failures.
Only a payload with no displayable market data is recorded as failed.
"""
if dashboard_has_usable_data(dashboard):
return dashboard
meta = dashboard.get("meta") if isinstance(dashboard, dict) else None
notice = ""
if isinstance(meta, dict):
notice = str(meta.get("notice") or meta.get("display_notice") or "")
return {
"status": "failed",
"error": notice or "未获取到可用行情",
}
def official_catchup_due(today: str, snapshot: dict[str, object]) -> bool:
now = datetime.now().astimezone().time().replace(tzinfo=None)
if not (dt_time(15, 5) <= now < dt_time(22, 0)):
return False
meta = snapshot.get("meta") if isinstance(snapshot.get("meta"), dict) else {}
actual = str(meta.get("trade_date") or "").replace("-", "")
if (
actual == today
and meta.get("limit_data_source") != "derived"
and not meta.get("carried_forward")
):
return False
return True
+11 -12
View File
@@ -5,16 +5,7 @@ import time
from datetime import date
from backend.bootstrap.config import normalize_date
def _verified_dashboard_result(dashboard: dict[str, object]) -> dict[str, object]:
meta = dashboard.get("meta") or {}
if isinstance(meta, dict) and meta.get("carried_forward"):
return {
"status": "failed",
"error": str(meta.get("notice") or "未获取到所选日期的最新行情"),
}
return dashboard
from backend.jobs.refresh import official_catchup_due, verified_dashboard_result
class JobServiceMixin:
@@ -36,7 +27,7 @@ class JobServiceMixin:
started = self.jobs.submit(
"market.refresh",
key,
lambda: _verified_dashboard_result(self.sync_dashboard(normalized)),
lambda: verified_dashboard_result(self.sync_dashboard(normalized)),
{"trade_date": normalized, "trigger": "administrator"},
)
return {"started": started, "job_key": key if started else ""}
@@ -54,7 +45,15 @@ class JobServiceMixin:
self.jobs.submit(
"market.refresh",
f"realtime:{today}:{bucket}",
lambda: self.sync_dashboard(today),
lambda: verified_dashboard_result(self.sync_dashboard(today)),
{"trade_date": today, "trigger": "realtime-poll"},
)
elif official_catchup_due(today, snapshot):
bucket = int(time.time() // 300)
self.jobs.submit(
"market.refresh",
f"catchup:{today}:{bucket}",
lambda: verified_dashboard_result(self.sync_dashboard(today)),
{"trade_date": today, "trigger": "official-catchup"},
)
self._schedule_automatic_screeners(today, snapshot)
+18 -12
View File
@@ -330,6 +330,7 @@
"system_service": "backend/features/system/service.py",
"account_bridge": "backend/features/accounts/application.py",
"job_lifecycle": "backend/jobs/service.py",
"job_refresh_status": "backend/jobs/refresh.py",
"feature_routes": "backend/features/*/routes.py"
},
"numeric_normalization": [
@@ -472,8 +473,8 @@
},
{
"path": "frontend/shared/shell.css",
"bytes": 63659,
"lines": 3763
"bytes": 63733,
"lines": 3767
},
{
"path": "backend/features/heaven/engine.py",
@@ -560,6 +561,11 @@
"bytes": 14743,
"lines": 342
},
{
"path": "frontend/shared/dashboard.js",
"bytes": 14740,
"lines": 316
},
{
"path": "frontend/shared/admin.js",
"bytes": 14410,
@@ -575,11 +581,6 @@
"bytes": 13219,
"lines": 289
},
{
"path": "frontend/shared/dashboard.js",
"bytes": 12894,
"lines": 274
},
{
"path": "backend/features/market/insights_auction_data.py",
"bytes": 12829,
@@ -785,16 +786,16 @@
"bytes": 2514,
"lines": 63
},
{
"path": "backend/jobs/service.py",
"bytes": 2337,
"lines": 59
},
{
"path": "backend/features/mentor/routes.py",
"bytes": 2299,
"lines": 57
},
{
"path": "backend/jobs/service.py",
"bytes": 2219,
"lines": 60
},
{
"path": "backend/features/screener/regime.py",
"bytes": 2202,
@@ -830,6 +831,11 @@
"bytes": 1791,
"lines": 46
},
{
"path": "backend/jobs/refresh.py",
"bytes": 1728,
"lines": 46
},
{
"path": "backend/features/alerts/routes.py",
"bytes": 1687,
+1 -1
View File
@@ -32,4 +32,4 @@
- 旧文档不能删:被替代的旧文档开头要加一行「⚠️ 本文档已过时,仅留档备查,请勿删除」,再写新版。
- 用中文大白话写,专业词要带通俗解释,让不懂代码的人也能看懂。
- 「问天」板块是冻结区,任何改动都不许碰;写文档时别误导后来人去改它
- 「问天」不是永久冻结区:此前只冻结过界面视觉方案,现已解冻。问天可纳入后续数据与功能迁移,不要再写成“永远不碰”
+1
View File
@@ -7,6 +7,7 @@
| 任务 | 说明 | 状态 |
|---|---|---|
| 全站视觉统一改造收尾 | 主线。17 个阶段已完成,正在最终验收、代码合并 | 收尾中 |
| 行情刷新误报与旧数据提示 | HEL-412:高级接口未到齐不再记整次失败;今日正式数据晚到时提示当前展示日期 | 施工中 |
| 手机端独立重新设计 | 先出视觉/交互规范和技术架构方案,等老板确认后再施工 | 方案送审中 |
## 已做完
+2 -2
View File
@@ -29,11 +29,11 @@
- **智能工具类(3 个)**:智能选股、问师、问天。
- **个人类(1 个)**:我的复盘。
其中「问天」是冻结区(见下面的硬规矩)
其中「问天」此前只在全站视觉改造阶段冻结过界面方案,现已解冻;问天可以纳入后续数据与功能迁移,但不等于本阶段要重做视觉
## 几条硬规矩(不能破坏的边界)
- 「问天」板块**冻结区**,任何改动都不许碰它
- 「问天」板块**不是永久冻结区**:此前冻结的是界面视觉方案,现已解冻。问天现有功能与界面不要破坏;后续数据与功能迁移可以纳入,不主动重做视觉
- **不用假数据冒充真行情**;数据缺失就明说“没有/不可用”,不能编。
- **每个用户自己的数据互相隔离**(自选、复盘、对话、问天历史等),看不到别人的。
- **计算由程序确定性完成**(情绪周期、智能选股、问天排盘等),AI 大模型(LLM,就是会聊天的那个 AI)只负责解释或编译自然语言条件,不能改计算结果。
+22
View File
@@ -770,11 +770,33 @@
scroll.classList.add("m-motion-fade-in");
}
function dashboardFreshnessNotice() {
const meta = (state.dashboard && state.dashboard.meta) || {};
if (meta.display_notice) return String(meta.display_notice);
const requested = String(meta.requested_date || "").replace(/-/g, "");
const actual = String(meta.trade_date || "").replace(/-/g, "");
const compact = actual;
const shown = /^\d{8}$/.test(compact)
? (Number(compact.slice(4, 6)) + " 月 " + Number(compact.slice(6, 8)) + " 日")
: "";
if (meta.data_status === "preparing" || (meta.carried_forward && actual && requested && actual !== requested)) {
return shown ? ("今日数据正在准备,当前展示 " + shown) : "今日数据正在准备,当前展示最近可用数据";
}
if (meta.data_status === "partial" || meta.limit_data_source === "derived") {
return meta.notice || "部分正式数据尚未到齐,当前展示日线推算结果";
}
return "";
}
function renderTopArea(key) {
const page = document.querySelector(".m-page");
if (!page) return;
let top = page.querySelector(".m-top");
let html = buildStrip();
const freshness = dashboardFreshnessNotice();
if (freshness) {
html = '<div class="m-phase-notice"><strong>' + escapeHtml(freshness) + "</strong></div>" + html;
}
if (key === "market/performance") html += performanceConclusion();
if (!top) {
top = document.createElement("div");
+49 -7
View File
@@ -66,11 +66,13 @@ async function startAdminRefresh() {
const requestedCompact = requestedDate.replaceAll("-", "");
const actualCompact = actualDate.replaceAll("-", "");
const updated = formatTimestamp(meta.updated_at);
if (actualCompact !== requestedCompact || meta.carried_forward) {
const reason = meta.notice ? `${meta.notice}` : "";
setAdminRefreshStatus("warning", `刷新已完成,但没有获取到 ${requestedDate} 的最新行情;当前仍是 ${actualDate || "未知日期"}${reason}`, "triangle-alert");
showToast("刷新完成,但未获取到所选日期的最新行情");
} else if (meta.notice) {
const freshness = dashboardFreshnessMessage(meta);
if (freshness || actualCompact !== requestedCompact || meta.carried_forward || meta.limit_data_source === "derived") {
setAdminRefreshStatus("warning", freshness || `部分正式数据尚未到齐,当前展示 ${actualDate || "最近可用数据"}`, "triangle-alert");
setStatus(freshness || "部分正式数据尚未到齐,当前展示最近可用数据");
return;
}
if (meta.notice) {
setAdminRefreshStatus("warning", `已刷新到 ${actualDate}${updated}),但数据源提示:${meta.notice}`, "triangle-alert");
showToast(`已刷新到 ${actualDate},请留意数据源提示`);
} else {
@@ -105,6 +107,37 @@ async function waitForAdminRefresh(jobKey) {
throw new Error("刷新等待超时,请稍后重试");
}
let dashboardCatchupTimer = 0;
function chineseMonthDay(value) {
const compact = String(value || "").replaceAll("-", "").replaceAll("/", "");
if (!/^\d{8}/.test(compact)) return "";
return `${Number(compact.slice(4, 6))}${Number(compact.slice(6, 8))}`;
}
function dashboardFreshnessMessage(meta = {}) {
if (meta.display_notice) return String(meta.display_notice);
const requested = String(meta.requested_date || "").replaceAll("-", "");
const actual = String(meta.trade_date || "").replaceAll("-", "");
const shown = chineseMonthDay(actual);
if (meta.data_status === "preparing" || (meta.carried_forward && actual && requested && actual !== requested)) {
return shown ? `今日数据正在准备,当前展示 ${shown}` : "今日数据正在准备,当前展示最近可用数据";
}
if (meta.data_status === "partial" || meta.limit_data_source === "derived") {
return meta.notice || "部分正式数据尚未到齐,当前展示日线推算结果";
}
return "";
}
function scheduleDashboardCatchup(meta = {}) {
window.clearTimeout(dashboardCatchupTimer);
const status = String(meta.data_status || "");
if (status !== "preparing" && status !== "partial") return;
dashboardCatchupTimer = window.setTimeout(() => {
loadDashboard(false, true, false);
}, 60000);
}
function applyDashboard(payload, background = false) {
state.dashboard = payload;
const selectedDate = payload.meta.requested_date || payload.meta.trade_date;
@@ -112,7 +145,11 @@ function applyDashboard(payload, background = false) {
document.querySelector("#qiObservationDate").value = selectedDate;
document.querySelector("#journalDate").value = selectedDate;
renderDashboard();
setStatus(`${dashboardSourceLabel(payload.meta)} · 数据已更新`);
const freshness = dashboardFreshnessMessage(payload.meta || {});
setStatus(freshness || `${dashboardSourceLabel(payload.meta)} · 数据已更新`);
const updatedAt = document.querySelector("#updatedAt");
if (updatedAt) updatedAt.dataset.tone = freshness ? "warning" : "ok";
scheduleDashboardCatchup(payload.meta || {});
if (!background) {
if (state.activeView === "dragonView") loadDragonTiger();
if (state.activeView === "screenerView") loadScreenerSetup();
@@ -180,7 +217,12 @@ function renderDashboard() {
}
}
updateSentimentGauge(overview.sentiment_score);
setText("updatedAt", `${dashboardSourceLabel(meta)} · 更新 ${formatTimestamp(meta.updated_at)}`);
const freshness = dashboardFreshnessMessage(meta);
setText("updatedAt", freshness
? freshness
: `${dashboardSourceLabel(meta)} · 更新 ${formatTimestamp(meta.updated_at)}`);
const updatedAt = document.querySelector("#updatedAt");
if (updatedAt) updatedAt.dataset.tone = freshness ? "warning" : "ok";
renderLimitTable();
renderLadderMini(ladders || []);
+4
View File
@@ -921,6 +921,10 @@ body.sidebar-collapsed .app-main {
text-align: right;
}
.status-bar #updatedAt[data-tone="warning"] {
color: var(--warning);
}
.status-bar .risk-note {
display: block;
+207 -7
View File
@@ -1,23 +1,223 @@
from __future__ import annotations
import copy
import threading
import unittest
from datetime import date, datetime, timedelta, timezone
from pathlib import Path
from backend.jobs.service import _verified_dashboard_result
from backend.features.market.service import MarketServiceMixin
from backend.jobs.refresh import (
dashboard_has_usable_data,
official_catchup_due,
verified_dashboard_result,
)
from backend.data.providers.tushare_transport import TushareError
class AdminRefreshStatusTests(unittest.TestCase):
def test_carried_snapshot_is_reported_as_failed_job(self):
result = _verified_dashboard_result(
{"meta": {"carried_forward": True, "notice": "官方涨跌停数据尚未返回"}}
def test_carried_snapshot_is_usable_not_failed_job(self):
result = verified_dashboard_result(
{
"meta": {
"trade_date": "2026-09-01",
"requested_date": "2026-09-02",
"carried_forward": True,
"notice": "今日数据正在准备,当前展示 9 月 1 日",
"data_status": "preparing",
},
"overview": {"limit_up_count": 12},
}
)
self.assertEqual(result["status"], "failed")
self.assertEqual(result["error"], "官方涨跌停数据尚未返回")
self.assertNotEqual(result.get("status"), "failed")
self.assertEqual(result["meta"]["data_status"], "preparing")
self.assertTrue(dashboard_has_usable_data(result))
def test_derived_limit_snapshot_is_usable_not_failed_job(self):
dashboard = {
"meta": {
"trade_date": "2026-09-02",
"limit_data_source": "derived",
"notice": "涨跌停高级接口当日数据尚未更新,已使用日线数据推算。",
"data_status": "partial",
},
"overview": {"limit_up_count": 8},
}
self.assertIs(verified_dashboard_result(dashboard), dashboard)
def test_current_snapshot_is_reported_as_successful_job(self):
dashboard = {"meta": {"trade_date": "2026-08-28", "carried_forward": False}}
self.assertIs(_verified_dashboard_result(dashboard), dashboard)
self.assertIs(verified_dashboard_result(dashboard), dashboard)
def test_empty_payload_is_still_failed(self):
result = verified_dashboard_result({"meta": {}, "overview": {}})
self.assertEqual(result["status"], "failed")
class FakeSyncDatabase:
def __init__(self, latest=None):
self.latest = latest
self.saved = []
self.finished = []
def start_sync(self, *_args, **_kwargs):
return 1
def save_snapshot(self, trade_date, source, payload):
self.saved.append((trade_date, source, copy.deepcopy(payload)))
def save_data_snapshot(self, *_args, **_kwargs):
return None
def finish_sync(self, *args, **kwargs):
self.finished.append((args, kwargs))
def get_latest_real_snapshot(self, *_args, **_kwargs):
return copy.deepcopy(self.latest)
def get_snapshot(self, *_args, **_kwargs):
return None
def get_data_snapshot(self, *_args, **_kwargs):
return None
def reason_overrides(self, *_args, **_kwargs):
return {}
class FakeDerivedClient:
def dashboard(self, trade_date: str):
return {
"meta": {
"trade_date": f"{trade_date[:4]}-{trade_date[4:6]}-{trade_date[6:8]}",
"limit_data_source": "derived",
"notice": "涨跌停高级接口当日数据尚未更新,已使用日线数据推算。",
"updated_at": datetime.now().astimezone().isoformat(timespec="seconds"),
},
"overview": {"limit_up_count": 3},
"limits": [{"code": "000001"}],
"broken": [],
"down_limits": [],
"yesterday_limits": [],
}
class FakeMissingDailyClient:
def dashboard(self, trade_date: str):
raise TushareError(f"No daily data returned for {trade_date}")
class SyncHarness(MarketServiceMixin):
def __init__(self, client, latest=None):
self.configured = True
self.sync_lock = threading.Lock()
self.database = FakeSyncDatabase(latest)
self._client = client
self.current_user_id = 1
def _tushare_client(self):
return self._client
def _enrich_dashboard_sentiment(self, dashboard, _trade_date):
return dashboard
def _apply_reason_overrides(self, dashboard):
return dashboard
class DashboardFreshnessTests(unittest.TestCase):
def test_derived_limits_are_kept_as_partial_success(self):
today = date.today().strftime("%Y%m%d")
harness = SyncHarness(FakeDerivedClient())
payload = harness.sync_dashboard(today)
meta = payload["meta"]
self.assertEqual(meta["limit_data_source"], "derived")
self.assertEqual(meta["data_status"], "partial")
self.assertFalse(meta.get("carried_forward"))
self.assertIn("日线数据推算", meta["display_notice"])
self.assertEqual(harness.database.finished[0][0][1], "success")
self.assertEqual(verified_dashboard_result(payload), payload)
def test_missing_official_data_keeps_previous_day_with_preparing_notice(self):
today = date.today()
previous = (today - timedelta(days=1)).strftime("%Y-%m-%d")
latest = {
"meta": {"trade_date": previous, "source": "tushare"},
"overview": {"limit_up_count": 20},
}
harness = SyncHarness(FakeMissingDailyClient(), latest)
payload = harness.sync_dashboard(today.strftime("%Y%m%d"))
meta = payload["meta"]
self.assertTrue(meta["carried_forward"])
self.assertEqual(meta["data_status"], "preparing")
self.assertIn("今日数据正在准备,当前展示", meta["display_notice"])
self.assertIn("", meta["display_notice"])
self.assertNotIn("No daily data", meta["display_notice"])
self.assertNotEqual(verified_dashboard_result(payload).get("status"), "failed")
def test_weekend_carry_is_not_labeled_as_preparing(self):
snapshot = {
"meta": {"trade_date": "2026-07-24", "source": "tushare", "updated_at": "2026-07-24T15:00:00+08:00"},
"overview": {"limit_up_count": 1},
}
harness = SyncHarness(FakeMissingDailyClient())
carried = harness._carry_dashboard(snapshot, "20260725", "非交易日沿用最近交易日收盘行情")
self.assertEqual(carried["meta"]["data_status"], "carried")
self.assertIn("非交易日", carried["meta"]["display_notice"])
def test_stale_derived_snapshot_is_retried(self):
today = date.today().strftime("%Y%m%d")
old = datetime.now(timezone.utc) - timedelta(minutes=5)
snapshot = {
"meta": {
"source": "tushare",
"trade_date": f"{today[:4]}-{today[4:6]}-{today[6:8]}",
"limit_data_source": "derived",
"updated_at": old.isoformat(),
},
"overview": {"limit_up_count": 1},
}
harness = SyncHarness(FakeDerivedClient())
harness.database.get_snapshot = lambda *_args, **_kwargs: copy.deepcopy(snapshot)
payload = harness.get_dashboard(today)
self.assertEqual(payload["meta"]["data_status"], "partial")
self.assertTrue(harness.database.saved)
def test_official_catchup_skips_complete_today_snapshot(self):
today = date.today().strftime("%Y%m%d")
iso = f"{today[:4]}-{today[4:6]}-{today[6:8]}"
due = official_catchup_due(
today,
{"meta": {"trade_date": iso, "limit_data_source": "official"}},
)
derived_due = official_catchup_due(
today,
{"meta": {"trade_date": iso, "limit_data_source": "derived"}},
)
now = datetime.now().astimezone().time().replace(tzinfo=None)
if datetime.strptime("15:05", "%H:%M").time() <= now < datetime.strptime("22:00", "%H:%M").time():
self.assertFalse(due)
self.assertTrue(derived_due)
else:
self.assertFalse(due)
self.assertFalse(derived_due)
class FrontendRefreshCopyTests(unittest.TestCase):
def test_dashboard_script_distinguishes_partial_from_failure(self):
script = (Path(__file__).resolve().parents[1] / "frontend" / "shared" / "dashboard.js").read_text(encoding="utf-8")
self.assertIn("今日数据正在准备,当前展示", script)
self.assertIn("部分正式数据尚未到齐", script)
self.assertIn('job.status === "failed"', script)
failed_block = script.split("if (job.status === \"failed\")", 1)[1].split("const query", 1)[0]
self.assertIn("后台刷新失败", failed_block)
success_block = script.split("const freshness = dashboardFreshnessMessage(meta);", 1)[1]
self.assertNotIn("后台刷新失败", success_block.split("} else {", 1)[0])
if __name__ == "__main__":
+34
View File
@@ -0,0 +1,34 @@
import logging
import unittest
from backend.bootstrap.runtime import configure_logging
class ConfigureLoggingTest(unittest.TestCase):
def setUp(self) -> None:
self._saved_handlers = logging.getLogger().handlers[:]
self._saved_level = logging.getLogger().level
logging.getLogger().handlers.clear()
def tearDown(self) -> None:
logging.getLogger().handlers[:] = self._saved_handlers
logging.getLogger().setLevel(self._saved_level)
def test_configures_root_logger_at_info(self) -> None:
configure_logging()
root = logging.getLogger()
self.assertTrue(root.handlers)
self.assertEqual(root.level, logging.INFO)
with self.assertLogs("xiaobai.datahub", level="INFO") as captured:
logging.getLogger("xiaobai.datahub").info("datahub shadow %s", {"dataset": "daily"})
self.assertIn("datahub shadow", captured.output[0])
def test_keeps_existing_configuration(self) -> None:
handler = logging.NullHandler()
logging.getLogger().addHandler(handler)
configure_logging()
self.assertEqual(logging.getLogger().handlers, [handler])
if __name__ == "__main__":
unittest.main()
+35 -2
View File
@@ -126,7 +126,7 @@ class DatahubBridgeTests(unittest.TestCase):
self.assertEqual(calendar[0]["is_open"], 1)
self.assertEqual(calendar_client.paths, [])
def test_fallback_on_down_401_timeout_empty_unpublished_and_stale(self) -> None:
def test_fallback_on_down_401_timeout_empty_unpublished_stale_and_incomplete(self) -> None:
cases = [
DatahubError("UNAVAILABLE", "down"),
DatahubError("UNAUTHORIZED", "401"),
@@ -134,6 +134,7 @@ class DatahubBridgeTests(unittest.TestCase):
DatahubError("EMPTY", "no rows"),
DatahubError("DATASET_NOT_PUBLISHED", "not ready"),
DatahubError("STALE", "old"),
DatahubError("INCOMPLETE", "truncated"),
]
for error in cases:
with self.subTest(error=error.code):
@@ -144,6 +145,16 @@ class DatahubBridgeTests(unittest.TestCase):
data=[dict(HUB_DAILY)],
meta={"stale": True, "staleness_seconds": 999999},
))
elif error.code == "INCOMPLETE":
client = FakeClient(response=DatahubResponse(
data=[dict(HUB_DAILY)],
meta={
"stale": False,
"staleness_seconds": 0,
"incomplete": True,
"coverage": {"complete": False, "missing_count": 80},
},
))
else:
client = FakeClient(error=error)
legacy = FakeLegacy([LEGACY_DAILY])
@@ -198,7 +209,8 @@ class DatahubBridgeTests(unittest.TestCase):
self.assertEqual(canonical["vol"], 100000.0)
self.assertEqual(canonical["amount"], 2000000.0)
def test_heaven_keeps_legacy_even_when_read_flag_is_on(self) -> None:
def test_heaven_keeps_legacy_on_first_batch_even_when_read_flag_is_on(self) -> None:
"""问天未永久冻结;首批只读接入仍走旧链路,后续迁移可以纳入。"""
self.assertTrue(looks_like_heaven("backend.features.heaven.market_context", "backend/features/heaven/market_context.py"))
self.assertFalse(looks_like_heaven("backend.features.market.service", "backend/features/market/service.py"))
client = FakeClient()
@@ -234,6 +246,27 @@ class DatahubBridgeTests(unittest.TestCase):
self.assertIsInstance(client, DatahubAwareTushareClient)
self.assertFalse(gateway.datahub.settings.any_enabled())
def test_stock_detail_range_query_is_not_silently_accepted_when_incomplete(self) -> None:
source = (ROOT / "backend" / "data" / "providers" / "tushare_stocks.py").read_text(encoding="utf-8")
self.assertIn('"daily"', source)
self.assertIn("start_date", source)
self.assertIn("end_date", source)
client = FakeClient(
response=DatahubResponse(
data=[dict(HUB_DAILY)],
meta={"stale": False, "staleness_seconds": 0, "incomplete": True, "coverage": {"complete": False, "missing_count": 89}},
)
)
legacy = FakeLegacy([LEGACY_DAILY])
wrapped = DatahubAwareTushareClient(legacy, DatahubBridge(flags(daily=(True, False)), client))
rows = wrapped.query(
"daily",
{"ts_code": "600000.SH", "start_date": "20240301", "end_date": "20240902"},
"ts_code,amount",
)
self.assertEqual(rows[0]["amount"], 2000.0)
self.assertEqual(len(legacy.calls), 1)
def test_features_do_not_import_datahub_client(self) -> None:
violations = []
for path in (ROOT / "backend" / "features").rglob("*.py"):
+2
View File
@@ -97,6 +97,7 @@ def code_hotspots() -> list[dict[str, Any]]:
"backend/features/system/service.py",
"backend/features/accounts/application.py",
"backend/jobs/service.py",
"backend/jobs/refresh.py",
"database.py",
"backend/features/screener/engine.py",
"backend/features/screener/catalog.py",
@@ -265,6 +266,7 @@ def build() -> dict[str, Any]:
"system_service": "backend/features/system/service.py",
"account_bridge": "backend/features/accounts/application.py",
"job_lifecycle": "backend/jobs/service.py",
"job_refresh_status": "backend/jobs/refresh.py",
"feature_routes": "backend/features/*/routes.py",
},
"numeric_normalization": [
+17
View File
@@ -63,6 +63,22 @@ python -m unittest discover -s tests -v
不调用真实 Tushare;用内存/临时库和假适配器。
## 历史回补
交易日历默认从 `20160101` 拉到今天后 30 天;盘前 `precheck` 与手动回补都走同一 UPSERT,可重复执行。
网站实际使用的指数(上证、深成、创业板、沪深300)按交易日增量发布,默认覆盖 260 个交易日(大于现有 90 天窗口,并覆盖智能选股基准回看)。已发布日期默认跳过。
```bash
cd xiaobai-datahub
python -m datahub history-backfill
# 可选:--calendar-start 20160101 --index-days 260 --force
```
管理后台也可手动跑 `history_backfill` 任务,或 `POST /admin/api/backfill``dataset=history`、确认词 `history:full`
区间接口在 `meta.coverage` / `meta.incomplete` 标明覆盖是否完整;网站只读接入把不完整区间视为不可用并回旧链路。个股日 K 的 90 天区间查询依赖已核实,本阶段不回补全市场历史。
## 备份
每日 00:40 任务把 `datahub.db` 备份到 `data/backups/`(保留 14 份)。也可手动:
@@ -74,5 +90,6 @@ python -c "from pathlib import Path; from datahub.db import HubDB; HubDB(Path('d
## 安全
- 密钥只以 `configured / 末4位 / 更新时间` 出现在后台,不进日志、不进 `/v1`
- HTTP 解析失败只记录“请求不是合法 JSON”,不把请求正文、密码或 Token 写入容器日志
- 回滚、补数需重新输入密码 + 确认词
- 容器非 rootuid 10002)、read_only、cap_drop ALL
@@ -11,5 +11,7 @@
"publication_generations": 3,
"tushare_rate_per_minute": 300,
"list_limit_default": 5000,
"list_limit_max": 5000
"list_limit_max": 5000,
"calendar_start": "20160101",
"index_history_trading_days": 260
}
+4
View File
@@ -0,0 +1,4 @@
from datahub.cli import main
if __name__ == "__main__":
raise SystemExit(main())
+7 -2
View File
@@ -44,7 +44,10 @@ DATASET_API = {
"auction": "stk_auction",
}
DEFAULT_INDEX_CODES = ("000001.SH", "399001.SZ", "399006.SZ", "000300.SH")
# Website actual index usage: market cards / 90-day charts (SH/SZ/CYB) plus
# screener 沪深300 benchmark (lookback up to 260 trading days).
WEBSITE_INDEX_CODES = ("000001.SH", "399001.SZ", "399006.SZ", "000300.SH")
DEFAULT_INDEX_CODES = WEBSITE_INDEX_CODES
class TushareAdapter(MarketAdapter):
@@ -145,7 +148,9 @@ class TushareAdapter(MarketAdapter):
try:
with urllib.request.urlopen(request, timeout=self.timeout) as response:
result = json.loads(response.read().decode("utf-8"))
except (urllib.error.URLError, TimeoutError, json.JSONDecodeError) as exc:
except json.JSONDecodeError:
raise AdapterError("Tushare returned invalid json") from None
except (urllib.error.URLError, TimeoutError) as exc:
raise AdapterError(f"Tushare request failed: {exc}") from exc
if result.get("code") != 0:
raise AdapterError(result.get("msg") or "Tushare returned an unknown error")
+11 -5
View File
@@ -88,6 +88,7 @@ class AdminAPI:
{"id": "precheck", "at": "08:45", "title": "盘前预检"},
{"id": "eod_a", "at": "15:05", "title": "盘后批 A daily/valuation/moneyflow/auction"},
{"id": "eod_b", "at": "15:10", "title": "盘后批 B index_daily"},
{"id": "history_backfill", "at": "manual", "title": "回补历史日历与指数日 K"},
{"id": "cleanup", "at": "00:30", "title": "清理 staging / 日志"},
{"id": "backup", "at": "00:40", "title": "SQLite 备份"},
],
@@ -130,12 +131,17 @@ class AdminAPI:
return result
def backfill(self, dataset: str, trade_date: str, password: str, confirm: str, actor: str) -> dict[str, Any]:
self._dangerous(password, confirm, f"{dataset}:{trade_date}")
if dataset == "reference":
result = self.pipeline.ingest_reference(trade_date)
day = yyyymmdd(trade_date or now_shanghai())
if dataset == "history":
self._dangerous(password, confirm, "history:full")
result = self.pipeline.backfill_history(day)
else:
result = self.pipeline.run_dataset(dataset, trade_date)
self.pipeline.audit(actor, "backfill", f"{dataset}:{trade_date}", json.dumps({"ok": True}))
self._dangerous(password, confirm, f"{dataset}:{day}")
if dataset == "reference":
result = self.pipeline.ingest_reference(day)
else:
result = self.pipeline.run_dataset(dataset, day)
self.pipeline.audit(actor, "backfill", f"{dataset}:{day}", json.dumps({"ok": True}))
return result
def _dangerous(self, password: str, confirm: str, expected: str) -> None:
+38
View File
@@ -0,0 +1,38 @@
"""Command-line entry for one-shot datahub operations."""
from __future__ import annotations
import argparse
import json
import sys
from datahub.hub import build_hub
from datahub.settings import load_settings
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(description="xiaobai-datahub CLI")
sub = parser.add_subparsers(dest="command", required=True)
history = sub.add_parser("history-backfill", help="回补 2016 年起交易日历和网站所用指数日 K")
history.add_argument("--calendar-start", default=None, help="日历起点,默认配置 calendar_start")
history.add_argument("--index-days", type=int, default=None, help="指数回补交易日数量,默认 260")
history.add_argument("--force", action="store_true", help="覆盖已发布的指数日期")
args = parser.parse_args(argv)
settings = load_settings()
hub = build_hub(settings)
if args.command == "history-backfill":
result = hub.pipeline.backfill_history(
calendar_start=args.calendar_start,
index_days=args.index_days,
force=args.force,
)
json.dump(result, sys.stdout, ensure_ascii=False, indent=2, default=str)
sys.stdout.write("\n")
return 0 if result.get("ok") else 1
parser.error(f"unknown command: {args.command}")
return 2
if __name__ == "__main__":
raise SystemExit(main())
+130
View File
@@ -0,0 +1,130 @@
from __future__ import annotations
from typing import Any, Iterable
from datahub.db import HubDB
from datahub.timeutil import iter_yyyymmdd, yyyymmdd
MISSING_SAMPLE_LIMIT = 10
def coverage_payload(
*,
kind: str,
start: str,
end: str,
expected: Iterable[str],
available: Iterable[str],
extra: dict[str, Any] | None = None,
) -> dict[str, Any]:
start = yyyymmdd(start)
end = yyyymmdd(end)
expected_list = sorted({yyyymmdd(item) for item in expected if item})
available_set = {yyyymmdd(item) for item in available if item}
missing = [item for item in expected_list if item not in available_set]
payload: dict[str, Any] = {
"kind": kind,
"complete": not missing,
"requested_from": start,
"requested_to": end,
"available_from": min(available_set) if available_set else None,
"available_to": max(available_set) if available_set else None,
"expected_count": len(expected_list),
"available_count": len(available_set),
"missing_count": len(missing),
"missing_sample": missing[:MISSING_SAMPLE_LIMIT],
}
if extra:
payload.update(extra)
return payload
def calendar_coverage(db: HubDB, start: str, end: str, exchange: str = "SSE") -> dict[str, Any]:
start = yyyymmdd(start)
end = yyyymmdd(end)
expected = list(iter_yyyymmdd(start, end))
rows = db.fetchall(
"SELECT cal_date FROM trade_calendar WHERE exchange = ? AND cal_date >= ? AND cal_date <= ?",
(exchange, start, end),
)
return coverage_payload(
kind="calendar",
start=start,
end=end,
expected=expected,
available=(row["cal_date"] for row in rows),
extra={"exchange": exchange},
)
def published_range_coverage(
db: HubDB,
dataset: str,
start: str,
end: str,
ts_code: str = "",
table: str = "",
) -> dict[str, Any]:
start = yyyymmdd(start)
end = yyyymmdd(end)
calendar = calendar_coverage(db, start, end)
open_rows = db.fetchall(
"""
SELECT cal_date FROM trade_calendar
WHERE exchange = 'SSE' AND is_open = 1 AND cal_date >= ? AND cal_date <= ?
ORDER BY cal_date
""",
(start, end),
)
expected_open = [row["cal_date"] for row in open_rows]
pubs = db.fetchall(
"""
SELECT trade_date, active_batch FROM publications
WHERE dataset = ? AND trade_date >= ? AND trade_date <= ?
ORDER BY trade_date
""",
(dataset, start, end),
)
published_dates = [row["trade_date"] for row in pubs]
available = list(published_dates)
extra: dict[str, Any] = {
"dataset": dataset,
"calendar_complete": calendar["complete"],
"calendar_missing_count": calendar["missing_count"],
}
if ts_code and table and pubs:
present_code: list[str] = []
for pub in pubs:
hit = db.fetchone(
f"SELECT 1 AS ok FROM {table} WHERE trade_date = ? AND batch_id = ? AND ts_code = ? LIMIT 1",
(pub["trade_date"], pub["active_batch"], ts_code),
)
if hit:
present_code.append(pub["trade_date"])
available = present_code
extra["code"] = ts_code
payload = coverage_payload(
kind="published_range",
start=start,
end=end,
expected=expected_open,
available=available,
extra=extra,
)
if not calendar["complete"]:
payload["complete"] = False
payload["calendar_missing_sample"] = calendar["missing_sample"]
return payload
def point_coverage(trade_date: str, dataset: str = "") -> dict[str, Any]:
day = yyyymmdd(trade_date)
payload = coverage_payload(
kind="point",
start=day,
end=day,
expected=[day],
available=[day],
extra={"dataset": dataset} if dataset else None,
)
return payload
+9 -1
View File
@@ -190,7 +190,15 @@ class HubRequestHandler(BaseHTTPRequestHandler):
return {}
if length <= 0 or length > 65536:
raise ValueError("请求内容为空或过大")
return json.loads(self.rfile.read(length).decode("utf-8"))
raw = self.rfile.read(length)
try:
payload = json.loads(raw.decode("utf-8"))
except (UnicodeDecodeError, json.JSONDecodeError):
LOGGER.warning("invalid json request body")
raise ValueError("请求不是合法 JSON") from None
if not isinstance(payload, dict):
raise ValueError("请求不是合法 JSON")
return payload
def _cookie_value(self, name: str) -> str:
cookie = SimpleCookie()
+28 -2
View File
@@ -2,7 +2,9 @@ from __future__ import annotations
import json
import logging
import re
import sys
import traceback
from typing import Any
from datahub.timeutil import isoformat
@@ -11,6 +13,13 @@ _SECRET_KEYS = (
"token", "password", "secret", "key", "authorization", "credential",
"tushare_token", "datahub_token", "encryption_key", "cookie",
)
_SECRET_JSON = re.compile(
r'(?i)("(?:' + "|".join(re.escape(key) for key in _SECRET_KEYS) + r')"\s*:\s*")([^"\\]*(?:\\.[^"\\]*)*)(")'
)
def redact_log_text(text: str) -> str:
return _SECRET_JSON.sub(r"\1***\3", str(text))
def _redact(value: Any, key: str = "") -> Any:
@@ -21,22 +30,39 @@ def _redact(value: Any, key: str = "") -> Any:
return {str(item_key): _redact(item_value, str(item_key)) for item_key, item_value in value.items()}
if isinstance(value, list):
return [_redact(item) for item in value]
if isinstance(value, str):
return redact_log_text(value)
return value
def _safe_exc_text(exc_info: tuple[Any, Any, Any]) -> str:
exc = exc_info[1]
if isinstance(exc, json.JSONDecodeError):
return f"JSONDecodeError: invalid json at position {exc.pos}"
cause = getattr(exc, "__cause__", None)
if isinstance(cause, json.JSONDecodeError):
return f"{type(exc).__name__}: invalid json in request"
text = "".join(traceback.format_exception(*exc_info))
if isinstance(cause, json.JSONDecodeError) and cause.doc:
text = text.replace(cause.doc, "")
if isinstance(exc, json.JSONDecodeError) and exc.doc:
text = text.replace(exc.doc, "")
return redact_log_text(text)
class JsonFormatter(logging.Formatter):
def format(self, record: logging.LogRecord) -> str:
payload: dict[str, Any] = {
"ts": isoformat(),
"level": record.levelname,
"logger": record.name,
"message": record.getMessage(),
"message": redact_log_text(record.getMessage()),
}
extra = getattr(record, "hub", None)
if isinstance(extra, dict):
payload.update(_redact(extra))
if record.exc_info:
payload["exc"] = self.formatException(record.exc_info)
payload["exc"] = _safe_exc_text(record.exc_info)
return json.dumps(payload, ensure_ascii=False, default=str)
+162 -7
View File
@@ -7,7 +7,7 @@ from datetime import timedelta
from typing import Any
from datahub.adapters.base import AdapterError
from datahub.adapters.tushare import DEFAULT_INDEX_CODES, TushareAdapter
from datahub.adapters.tushare import DEFAULT_INDEX_CODES, WEBSITE_INDEX_CODES, TushareAdapter
from datahub.db import DATASET_TABLES, HubDB
from datahub.governance.circuit import CircuitBreaker
from datahub.governance.ratelimit import TokenBucket
@@ -141,11 +141,22 @@ class Pipeline:
seq = int((row or {}).get("n") or 0) + 1
return f"{trade_date}-{dataset}-{seq:03d}"
def ingest_reference(self, trade_date: str | None = None) -> dict[str, Any]:
"""Refresh trade calendar (window) and stock master. Not versioned by batch."""
def ingest_reference(
self,
trade_date: str | None = None,
start: str | None = None,
end: str | None = None,
) -> dict[str, Any]:
"""Refresh trade calendar and stock master. Not versioned by batch.
Calendar defaults to 2016-01-01 through today+30 so a 5-year website
query is not silently truncated. UPSERT makes repeats safe.
"""
day = yyyymmdd(trade_date or self.clock())
start = add_days(day, -400)
end = add_days(day, 30)
start = yyyymmdd(start or self.settings.calendar_start)
end = yyyymmdd(end or add_days(day, 30))
if start > end:
start, end = end, start
calendar = self.adapter.normalize(
"calendar",
self._guarded_fetch("calendar", {"exchange": "SSE", "start_date": start, "end_date": end}),
@@ -180,9 +191,150 @@ class Pipeline:
row.get("list_date"), fetched_at,
),
)
return {"calendar": len(calendar), "stocks": len(stocks), "trade_date": day}
return {
"calendar": len(calendar),
"stocks": len(stocks),
"trade_date": day,
"calendar_from": start,
"calendar_to": end,
}
def run_dataset(self, dataset: str, trade_date: str, attempts: int | None = None) -> dict[str, Any]:
def open_trade_dates(self, end: str, limit: int) -> list[str]:
end = yyyymmdd(end)
rows = self.db.fetchall(
"""
SELECT cal_date FROM trade_calendar
WHERE exchange = 'SSE' AND is_open = 1 AND cal_date <= ?
ORDER BY cal_date DESC
LIMIT ?
""",
(end, max(1, int(limit))),
)
return sorted(str(row["cal_date"]) for row in rows)
def backfill_history(
self,
trade_date: str | None = None,
calendar_start: str | None = None,
index_days: int | None = None,
codes: tuple[str, ...] | None = None,
force: bool = False,
) -> dict[str, Any]:
"""Idempotent calendar + website-index history backfill."""
day = yyyymmdd(trade_date or self.clock())
calendar = self.ingest_reference(day, start=calendar_start)
index = self.backfill_index_history(
end_date=day,
trading_days=index_days,
codes=codes,
force=force,
)
return {"calendar": calendar, "index_daily": index, "ok": bool(index.get("ok"))}
def backfill_index_history(
self,
end_date: str | None = None,
trading_days: int | None = None,
codes: tuple[str, ...] | None = None,
force: bool = False,
) -> dict[str, Any]:
"""Incrementally publish official index bars for website index codes.
One range fetch per code, then per-day publish. Already published dates
are skipped unless ``force``. Failures are recorded and do not roll back
successful days.
"""
end = yyyymmdd(end_date or self.clock())
limit = int(trading_days or self.settings.index_history_trading_days)
codes = tuple(codes or WEBSITE_INDEX_CODES)
open_dates = self.open_trade_dates(end, limit)
if not open_dates:
return {
"start": None,
"end": end,
"codes": list(codes),
"requested_days": 0,
"published": [],
"skipped": [],
"failed": [{"error": "calendar has no open dates on or before end"}],
"ok": False,
}
start = open_dates[0]
complete_dates = set() if force else self._index_dates_with_all_codes(start, end, codes)
targets = [day for day in open_dates if day not in complete_dates]
skipped = [day for day in open_dates if day in complete_dates]
by_date: dict[str, list[dict[str, Any]]] = {day: [] for day in targets}
failed: list[dict[str, Any]] = []
for ts_code in codes:
try:
raw = retry_call(
lambda code=ts_code: self._guarded_fetch(
"index_daily",
{"ts_code": code, "start_date": start, "end_date": end},
),
attempts=self.settings.max_publish_attempts,
base_delay=0.05,
sleeper=lambda _d: time.sleep(_d),
)
for row in self.adapter.normalize("index_daily", raw):
day = str(row.get("trade_date") or "")
if day in by_date:
by_date[day].append(row)
except Exception as exc:
failed.append({"ts_code": ts_code, "error": str(exc)})
published: list[dict[str, Any]] = []
for day in targets:
rows = by_date.get(day) or []
try:
result = self.run_dataset("index_daily", day, prepared_rows=rows)
published.append(
{
"trade_date": day,
"batch_id": result["batch_id"],
"rows": result["rows"],
"state": result["state"],
}
)
except Exception as exc:
failed.append({"trade_date": day, "error": str(exc), "rows": len(rows)})
return {
"start": start,
"end": end,
"codes": list(codes),
"requested_days": len(open_dates),
"published": published,
"skipped": skipped,
"failed": failed,
"ok": not failed,
}
def _index_dates_with_all_codes(self, start: str, end: str, codes: tuple[str, ...]) -> set[str]:
pubs = self.db.fetchall(
"""
SELECT trade_date, active_batch FROM publications
WHERE dataset = 'index_daily' AND trade_date >= ? AND trade_date <= ?
""",
(start, end),
)
needed = set(codes)
complete: set[str] = set()
for pub in pubs:
rows = self.db.fetchall(
"SELECT DISTINCT ts_code FROM eod_index_bars WHERE trade_date = ? AND batch_id = ?",
(pub["trade_date"], pub["active_batch"]),
)
have = {str(row["ts_code"]) for row in rows}
if needed <= have:
complete.add(str(pub["trade_date"]))
return complete
def run_dataset(
self,
dataset: str,
trade_date: str,
attempts: int | None = None,
prepared_rows: list[dict[str, Any]] | None = None,
) -> dict[str, Any]:
trade_date = yyyymmdd(trade_date)
batch_id = self.next_batch_id(dataset, trade_date)
max_attempts = attempts or self.settings.max_publish_attempts
@@ -190,12 +342,15 @@ class Pipeline:
rows: list[dict[str, Any]] = []
try:
self._set_batch(batch_id, dataset, trade_date, "fetching", 1)
if prepared_rows is None:
rows = retry_call(
lambda: self._fetch_dataset(dataset, trade_date),
attempts=max_attempts,
base_delay=0.05,
sleeper=lambda _d: None if attempts == 1 else time.sleep(_d),
)
else:
rows = list(prepared_rows)
self._stage(dataset, batch_id, rows)
self._set_batch(batch_id, dataset, trade_date, "staged", 1, rows_in=len(rows), rows_out=len(rows))
self._set_batch(batch_id, dataset, trade_date, "validating", 1)
+4
View File
@@ -37,6 +37,7 @@ class Scheduler:
"eod_b": self._eod_b,
"cleanup": self._cleanup,
"backup": self._backup,
"history_backfill": self._history_backfill,
}
self._stop = threading.Event()
self._thread: threading.Thread | None = None
@@ -125,6 +126,9 @@ class Scheduler:
def _eod_b(self, trade_date: str) -> dict[str, Any]:
return self.pipeline.run_eod_batch_b(trade_date)
def _history_backfill(self, trade_date: str) -> dict[str, Any]:
return self.pipeline.backfill_history(trade_date)
def _cleanup(self, trade_date: str) -> dict[str, Any]:
result = self.pipeline.cleanup()
if now_shanghai().weekday() == 6:
+22 -2
View File
@@ -6,6 +6,7 @@ from urllib.parse import parse_qs
from datahub import SCHEMA_VERSION
from datahub.codes import resolve_code
from datahub.coverage import calendar_coverage, point_coverage, published_range_coverage
from datahub.db import HubDB
from datahub.normalize import qfq_bar
from datahub.numbers import finite_number
@@ -129,7 +130,8 @@ class V1API:
}
for row in rows
]
return envelope(items, self._official_meta("calendar", end if items else start, source="tushare:trade_cal"))
meta = self._official_meta("calendar", end if items else start, source="tushare:trade_cal")
return envelope(items, attach_coverage(meta, calendar_coverage(self.db, start, end)))
def stocks(self, updated_since: str, q: dict[str, str]) -> dict[str, Any]:
limit, offset = self._page(q)
@@ -272,7 +274,7 @@ class V1API:
"staleness_seconds": 0,
"state": pub["state"],
}
return envelope(rows, meta)
return envelope(rows, attach_coverage(meta, point_coverage(start, dataset)))
# multi-day: walk published dates
pubs = self.db.fetchall(
"SELECT * FROM publications WHERE dataset = ? AND trade_date >= ? AND trade_date <= ? ORDER BY trade_date",
@@ -294,8 +296,17 @@ class V1API:
if adjust == "qfq" and dataset == "daily":
sliced = self._apply_qfq(sliced)
last = pubs[-1]
coverage = published_range_coverage(
self.db,
dataset,
start,
end,
ts_code=ts_code,
table=table,
)
return envelope(
sliced,
attach_coverage(
{
"tier": "official",
"trade_date": last["trade_date"],
@@ -305,6 +316,8 @@ class V1API:
"stale": False,
"staleness_seconds": 0,
},
coverage,
),
)
def _apply_qfq(self, rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
@@ -359,6 +372,13 @@ def add_default(days: int) -> str:
return (now_shanghai() + timedelta(days=days)).strftime("%Y%m%d")
def attach_coverage(meta: dict[str, Any], coverage: dict[str, Any]) -> dict[str, Any]:
merged = dict(meta)
merged["coverage"] = coverage
merged["incomplete"] = not bool(coverage.get("complete"))
return merged
def parse_query(raw: str) -> dict[str, list[str]]:
return parse_qs(raw, keep_blank_values=True)
+8
View File
@@ -48,6 +48,14 @@ class Settings:
def list_limit_max(self) -> int:
return int(self.quality.get("list_limit_max") or 5000)
@property
def calendar_start(self) -> str:
return str(self.quality.get("calendar_start") or "20160101")
@property
def index_history_trading_days(self) -> int:
return int(self.quality.get("index_history_trading_days") or 260)
def load_settings(
env: dict[str, str] | None = None,
+10
View File
@@ -58,6 +58,16 @@ def add_days(trade_date: str, days: int) -> str:
return (parse_trade_date(trade_date) + timedelta(days=days)).strftime("%Y%m%d")
def iter_yyyymmdd(start: str, end: str):
cursor = parse_trade_date(start)
last = parse_trade_date(end)
if cursor > last:
return
while cursor <= last:
yield cursor.strftime("%Y%m%d")
cursor += timedelta(days=1)
def utc_timestamp(value: Any) -> str:
if isinstance(value, datetime):
return isoformat(value)
+11 -1
View File
@@ -51,7 +51,17 @@ RAW = {
def fake_transport(api_name: str, params: dict, fields: str):
if api_name == "index_daily":
code = params.get("ts_code")
return [row for row in RAW["index_daily"] if row["ts_code"] == code]
rows = [row for row in RAW["index_daily"] if row["ts_code"] == code]
trade_date = str(params.get("trade_date") or "")
start = str(params.get("start_date") or "")
end = str(params.get("end_date") or "")
if trade_date:
rows = [row for row in rows if row["trade_date"] == trade_date]
if start:
rows = [row for row in rows if row["trade_date"] >= start]
if end:
rows = [row for row in rows if row["trade_date"] <= end]
return rows
if api_name == "trade_cal":
start = str(params.get("start_date") or "")
end = str(params.get("end_date") or "99999999")
+51
View File
@@ -1,17 +1,21 @@
from __future__ import annotations
import io
import json
import logging
import tempfile
import threading
import unittest
from http.server import ThreadingHTTPServer
from pathlib import Path
from urllib.error import HTTPError
from urllib.request import Request, urlopen
from datahub.adapters.tushare import TushareAdapter
from datahub.crypto import SecretVault
from datahub.httpapp import make_handler
from datahub.hub import Hub
from datahub.logutil import JsonFormatter
from datahub.settings import Settings
from tests.fixtures import fake_transport
@@ -91,6 +95,53 @@ class AdminTests(unittest.TestCase):
)
self.assertEqual(ctx.exception.code, 401)
def test_invalid_json_does_not_log_request_body_secrets(self) -> None:
secret = "SuperSecretPass1!"
token = "hub-token-should-not-leak"
raw = json.dumps({"password": secret, "token": token, "username": "hub_admin"}) + "{not-json"
stream = io.StringIO()
logger = logging.getLogger("datahub")
handler = logging.StreamHandler(stream)
handler.setFormatter(JsonFormatter())
logger.addHandler(handler)
previous_level = logger.level
logger.setLevel(logging.DEBUG)
try:
req = Request(
self.base + "/admin/api/login",
data=raw.encode("utf-8"),
headers={"Content-Type": "application/json"},
method="POST",
)
with self.assertRaises(HTTPError) as ctx:
urlopen(req, timeout=5)
body = ctx.exception.read().decode("utf-8")
self.assertEqual(ctx.exception.code, 400)
self.assertNotIn(secret, body)
self.assertNotIn(token, body)
blob = stream.getvalue() + body
self.assertNotIn(secret, blob)
self.assertNotIn(token, blob)
self.assertNotIn(raw, blob)
finally:
logger.removeHandler(handler)
logger.setLevel(previous_level)
def test_json_formatter_drops_decode_error_document(self) -> None:
secret = "ParseSecretTokenXYZ"
formatter = JsonFormatter()
logger = logging.getLogger("datahub.test")
record = logger.makeRecord(
"datahub.test", logging.ERROR, __file__, 1, "parse failed", (), None
)
try:
json.loads('{"password": "%s"}{' % secret)
except json.JSONDecodeError as exc:
record.exc_info = (type(exc), exc, exc.__traceback__)
blob = formatter.format(record)
self.assertNotIn(secret, blob)
self.assertIn("invalid json", blob)
if __name__ == "__main__":
unittest.main()
+3
View File
@@ -107,6 +107,9 @@ class ApiContractTests(unittest.TestCase):
self.assertIn("data", body)
self.assertIn("meta", body)
self.assertIn("tier", body["meta"])
if "calendar" in path or "bars" in path or "indexes" in path or "valuation" in path or "moneyflow" in path or "auction" in path:
self.assertIn("coverage", body["meta"])
self.assertIn("incomplete", body["meta"])
def test_qfq_matches_formula(self) -> None:
_, none = self._get(f"/v1/bars/daily?date={TRADE_DATE}&code=600000.SH&adjust=none", token=self.token)
@@ -0,0 +1,236 @@
from __future__ import annotations
import unittest
from datetime import date, timedelta
from datahub.coverage import calendar_coverage, point_coverage, published_range_coverage
from datahub.serving import V1API
from tests.fixtures import TRADE_DATE, fake_transport
from tests.test_pipeline import make_pipeline
def history_transport(open_dates: list[str], extra_closed: list[str] | None = None):
open_set = set(open_dates)
start = date(int(open_dates[0][:4]), int(open_dates[0][4:6]), int(open_dates[0][6:8]))
end = date(int(open_dates[-1][:4]), int(open_dates[-1][4:6]), int(open_dates[-1][6:8]))
calendar = []
cursor = start
while cursor <= end:
compact = cursor.strftime("%Y%m%d")
calendar.append(
{
"exchange": "SSE",
"cal_date": compact,
"is_open": 1 if compact in open_set else 0,
"pretrade_date": compact,
}
)
cursor += timedelta(days=1)
for day in extra_closed or []:
calendar.append(
{"exchange": "SSE", "cal_date": day, "is_open": 0, "pretrade_date": open_dates[0]}
)
index_codes = ("000001.SH", "399001.SZ", "399006.SZ", "000300.SH")
index_rows = []
for ts_code in index_codes:
for day in open_dates:
index_rows.append(
{
"ts_code": ts_code,
"trade_date": day,
"open": 100,
"high": 101,
"low": 99,
"close": 100.5,
"pct_chg": 0.1,
"vol": 10.0,
"amount": 20.0,
}
)
def transport(api_name, params, fields):
if api_name == "trade_cal":
start = str(params.get("start_date") or "")
end = str(params.get("end_date") or "99999999")
return [row for row in calendar if start <= row["cal_date"] <= end]
if api_name == "index_daily":
code = params.get("ts_code")
rows = [row for row in index_rows if row["ts_code"] == code]
trade_date = str(params.get("trade_date") or "")
start = str(params.get("start_date") or "")
end = str(params.get("end_date") or "")
if trade_date:
rows = [row for row in rows if row["trade_date"] == trade_date]
if start:
rows = [row for row in rows if row["trade_date"] >= start]
if end:
rows = [row for row in rows if row["trade_date"] <= end]
return rows
return fake_transport(api_name, params, fields)
return transport
def consecutive_open_days(end: str, count: int) -> list[str]:
cursor = date(int(end[:4]), int(end[4:6]), int(end[6:8]))
days: list[str] = []
while len(days) < count:
if cursor.weekday() < 5:
days.append(cursor.strftime("%Y%m%d"))
cursor -= timedelta(days=1)
return sorted(days)
class CoverageApiTests(unittest.TestCase):
def test_calendar_marks_holes_incomplete(self) -> None:
pipe, _db = make_pipeline()
pipe.ingest_reference(TRADE_DATE)
api = V1API(pipe.db, pipe, pipe.settings)
payload = api.handle("/v1/calendar", {"from": ["20240901"], "to": ["20240907"]})
self.assertTrue(payload["meta"]["incomplete"])
self.assertFalse(payload["meta"]["coverage"]["complete"])
self.assertGreater(payload["meta"]["coverage"]["missing_count"], 0)
self.assertIn("20240901", payload["meta"]["coverage"]["missing_sample"])
def test_calendar_complete_when_every_day_present(self) -> None:
pipe, _db = make_pipeline()
pipe.ingest_reference(TRADE_DATE)
api = V1API(pipe.db, pipe, pipe.settings)
payload = api.handle("/v1/calendar", {"from": ["20240902"], "to": ["20240903"]})
self.assertFalse(payload["meta"]["incomplete"])
self.assertTrue(payload["meta"]["coverage"]["complete"])
self.assertEqual(payload["meta"]["coverage"]["expected_count"], 2)
self.assertEqual(len(payload["data"]), 2)
def test_index_range_incomplete_without_history(self) -> None:
pipe, _db = make_pipeline()
pipe.ingest_reference(TRADE_DATE)
pipe.run_dataset("index_daily", TRADE_DATE)
api = V1API(pipe.db, pipe, pipe.settings)
payload = api.handle(
"/v1/indexes/bars",
{"from": ["20240902"], "to": ["20240903"], "code": ["000001.SH"]},
)
self.assertTrue(payload["meta"]["incomplete"])
self.assertFalse(payload["meta"]["coverage"]["complete"])
self.assertEqual(payload["meta"]["coverage"]["available_count"], 1)
self.assertIn("20240903", payload["meta"]["coverage"]["missing_sample"])
def test_index_point_query_stays_complete(self) -> None:
pipe, _db = make_pipeline()
pipe.ingest_reference(TRADE_DATE)
pipe.run_dataset("index_daily", TRADE_DATE)
api = V1API(pipe.db, pipe, pipe.settings)
payload = api.handle("/v1/indexes/bars", {"date": [TRADE_DATE], "code": ["000001.SH"]})
self.assertFalse(payload["meta"]["incomplete"])
self.assertTrue(payload["meta"]["coverage"]["complete"])
self.assertEqual(payload["meta"]["coverage"]["kind"], "point")
def test_daily_range_incomplete_without_stock_history(self) -> None:
pipe, _db = make_pipeline()
pipe.ingest_reference(TRADE_DATE)
pipe.run_dataset("daily", TRADE_DATE)
api = V1API(pipe.db, pipe, pipe.settings)
payload = api.handle(
"/v1/bars/daily",
{"from": ["20240902"], "to": ["20240903"], "code": ["600000.SH"]},
)
self.assertTrue(payload["meta"]["incomplete"])
self.assertFalse(payload["meta"]["coverage"]["complete"])
class HistoryBackfillTests(unittest.TestCase):
def test_index_history_is_idempotent_and_covers_requested_days(self) -> None:
open_dates = consecutive_open_days(TRADE_DATE, 5)
pipe, db = make_pipeline(quality={"index_history_trading_days": 5, "calendar_start": open_dates[0]})
pipe.adapter._transport = history_transport(open_dates)
first = pipe.backfill_history(TRADE_DATE, index_days=5)
self.assertTrue(first["ok"])
self.assertEqual(first["calendar"]["calendar_from"], open_dates[0])
self.assertEqual(first["index_daily"]["requested_days"], 5)
self.assertEqual(len(first["index_daily"]["published"]), 5)
self.assertEqual(first["index_daily"]["skipped"], [])
pubs = db.fetchall("SELECT trade_date FROM publications WHERE dataset='index_daily'")
self.assertEqual(sorted(row["trade_date"] for row in pubs), open_dates)
second = pipe.backfill_index_history(TRADE_DATE, trading_days=5)
self.assertTrue(second["ok"])
self.assertEqual(second["published"], [])
self.assertEqual(second["skipped"], open_dates)
api = V1API(db, pipe, pipe.settings)
payload = api.handle(
"/v1/indexes/bars",
{"from": [open_dates[0]], "to": [open_dates[-1]], "code": ["000001.SH"]},
)
self.assertFalse(payload["meta"]["incomplete"])
self.assertEqual(payload["meta"]["coverage"]["available_count"], 5)
self.assertEqual(len(payload["data"]), 5)
def test_index_history_retries_failed_dates_without_dropping_success(self) -> None:
open_dates = consecutive_open_days(TRADE_DATE, 3)
base = history_transport(open_dates)
def missing_cyb(api_name, params, fields):
if api_name == "index_daily" and params.get("ts_code") == "399006.SZ":
raise RuntimeError("upstream down")
return base(api_name, params, fields)
pipe, db = make_pipeline(quality={"index_history_trading_days": 3, "max_publish_attempts": 1})
pipe.adapter._transport = missing_cyb
first = pipe.backfill_history(TRADE_DATE, calendar_start=open_dates[0], index_days=3)
self.assertFalse(first["ok"])
self.assertTrue(any(item.get("ts_code") == "399006.SZ" for item in first["index_daily"]["failed"]))
published_first = {
row["trade_date"]
for row in db.fetchall("SELECT trade_date FROM publications WHERE dataset='index_daily'")
}
self.assertEqual(published_first, set(open_dates))
pipe.adapter._transport = base
retry = pipe.backfill_index_history(TRADE_DATE, trading_days=3)
self.assertTrue(retry["ok"])
self.assertEqual(len(retry["published"]), 3)
for day in open_dates:
rows = db.fetchall(
"""
SELECT DISTINCT ts_code FROM eod_index_bars
WHERE trade_date = ? AND batch_id = (
SELECT active_batch FROM publications
WHERE dataset='index_daily' AND trade_date = ?
)
""",
(day, day),
)
self.assertEqual({row["ts_code"] for row in rows}, {"000001.SH", "399001.SZ", "399006.SZ", "000300.SH"})
def test_prepared_rows_skip_upstream_fetch(self) -> None:
pipe, _db = make_pipeline()
pipe.ingest_reference(TRADE_DATE)
calls = {"n": 0}
original = pipe.adapter._transport
def counting(api_name, params, fields):
calls["n"] += 1
return original(api_name, params, fields)
pipe.adapter._transport = counting
rows = pipe.adapter.normalize("index_daily", original("index_daily", {"ts_code": "000001.SH", "trade_date": TRADE_DATE}, ""))
before = calls["n"]
result = pipe.run_dataset("index_daily", TRADE_DATE, prepared_rows=rows)
self.assertEqual(result["rows"], 1)
self.assertEqual(calls["n"], before)
def test_coverage_helpers_point_and_calendar(self) -> None:
pipe, db = make_pipeline()
pipe.ingest_reference(TRADE_DATE)
point = point_coverage(TRADE_DATE, "index_daily")
self.assertTrue(point["complete"])
cal = calendar_coverage(db, "20240902", "20240903")
self.assertTrue(cal["complete"])
pub = published_range_coverage(db, "index_daily", "20240902", "20240903")
self.assertFalse(pub["complete"])
if __name__ == "__main__":
unittest.main()
+1
View File
@@ -69,6 +69,7 @@ class PipelineTests(unittest.TestCase):
pipe, db = make_pipeline()
ref = pipe.ingest_reference(TRADE_DATE)
self.assertEqual(ref["stocks"], 2)
self.assertEqual(ref["calendar_from"], "20160101")
result = pipe.run_dataset("daily", TRADE_DATE)
self.assertEqual(result["state"], "published")
self.assertEqual(result["rows"], 2)