Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a836cda1b2 | ||
|
|
25ff6bbe06 | ||
|
|
5085cacf0d |
+3
-1
@@ -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.
|
||||
|
||||
@@ -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。
|
||||
- 股市有风险,入市需谨慎。投资决策及其后果由使用者本人承担。
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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,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:
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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
@@ -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)
|
||||
|
||||
@@ -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
@@ -32,4 +32,4 @@
|
||||
|
||||
- 旧文档不能删:被替代的旧文档开头要加一行「⚠️ 本文档已过时,仅留档备查,请勿删除」,再写新版。
|
||||
- 用中文大白话写,专业词要带通俗解释,让不懂代码的人也能看懂。
|
||||
- 「问天」板块是冻结区,任何改动都不许碰;写文档时别误导后来人去改它。
|
||||
- 「问天」不是永久冻结区:此前只冻结过界面视觉方案,现已解冻。问天可纳入后续数据与功能迁移,不要再写成“永远不碰”。
|
||||
|
||||
@@ -7,6 +7,7 @@
|
||||
| 任务 | 说明 | 状态 |
|
||||
|---|---|---|
|
||||
| 全站视觉统一改造收尾 | 主线。17 个阶段已完成,正在最终验收、代码合并 | 收尾中 |
|
||||
| 行情刷新误报与旧数据提示 | HEL-412:高级接口未到齐不再记整次失败;今日正式数据晚到时提示当前展示日期 | 施工中 |
|
||||
| 手机端独立重新设计 | 先出视觉/交互规范和技术架构方案,等老板确认后再施工 | 方案送审中 |
|
||||
|
||||
## 已做完
|
||||
|
||||
+2
-2
@@ -29,11 +29,11 @@
|
||||
- **智能工具类(3 个)**:智能选股、问师、问天。
|
||||
- **个人类(1 个)**:我的复盘。
|
||||
|
||||
其中「问天」是冻结区(见下面的硬规矩)。
|
||||
其中「问天」此前只在全站视觉改造阶段冻结过界面方案,现已解冻;问天可以纳入后续数据与功能迁移,但不等于本阶段要重做视觉。
|
||||
|
||||
## 几条硬规矩(不能破坏的边界)
|
||||
|
||||
- 「问天」板块是**冻结区**,任何改动都不许碰它。
|
||||
- 「问天」板块**不是永久冻结区**:此前冻结的是界面视觉方案,现已解冻。问天现有功能与界面不要破坏;后续数据与功能迁移可以纳入,不主动重做视觉。
|
||||
- **不用假数据冒充真行情**;数据缺失就明说“没有/不可用”,不能编。
|
||||
- **每个用户自己的数据互相隔离**(自选、复盘、对话、问天历史等),看不到别人的。
|
||||
- **计算由程序确定性完成**(情绪周期、智能选股、问天排盘等),AI 大模型(LLM,就是会聊天的那个 AI)只负责解释或编译自然语言条件,不能改计算结果。
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -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 || []);
|
||||
|
||||
@@ -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;
|
||||
|
||||
|
||||
@@ -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__":
|
||||
|
||||
@@ -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()
|
||||
@@ -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"):
|
||||
|
||||
@@ -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": [
|
||||
|
||||
@@ -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 写入容器日志
|
||||
- 回滚、补数需重新输入密码 + 确认词
|
||||
- 容器非 root(uid 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
|
||||
}
|
||||
|
||||
@@ -0,0 +1,4 @@
|
||||
from datahub.cli import main
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -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")
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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())
|
||||
@@ -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
|
||||
@@ -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()
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
|
||||
@@ -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)
|
||||
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),
|
||||
)
|
||||
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)
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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,17 +296,28 @@ 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,
|
||||
{
|
||||
"tier": "official",
|
||||
"trade_date": last["trade_date"],
|
||||
"published_at": last["published_at"],
|
||||
"source": source,
|
||||
"batch_id": last["active_batch"],
|
||||
"stale": False,
|
||||
"staleness_seconds": 0,
|
||||
},
|
||||
attach_coverage(
|
||||
{
|
||||
"tier": "official",
|
||||
"trade_date": last["trade_date"],
|
||||
"published_at": last["published_at"],
|
||||
"source": source,
|
||||
"batch_id": last["active_batch"],
|
||||
"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)
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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()
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user