Compare commits

..
Author SHA1 Message Date
abd4d22a67 HEL-543 返工: 补充 observability 运行时可关闭开关 (DATAHUB_OBSERVABILITY)
总工复核 🔴:安全边界要求新增观测功能必须可关闭、关闭后现有功能完全照旧,
但此前实现没有任何运行时开关。修复:

- settings.py: 新增 Settings.observability_enabled 字段,沿用既有
  DATAHUB_SCHEDULER 的环境变量模式,读取 DATAHUB_OBSERVABILITY
  (0/false/off 关闭,默认开启)。
- hub.py: Hub.__init__ 把 settings.observability_enabled 挂到
  self.db 上,让 pipeline/realtime_serve/steward/admin_api 已经
  在传的 db 参数直接带上开关,零额外改造。
- observability.py: 新增 is_enabled(db),缺失该属性时默认按启用处理
  (向后兼容裸 HubDB 用例/测试)。关闭时 observe() 变成纯
  透传(不计时、不分类、不碰数据库),record_call() 直接 no-op。
- admin_api.py: 4 个新只读端点关闭时返回明确的
  {"enabled": false, ...空结构} 而不是静默返回旧数据。
- 新增 10 个测试:开关默认值/环境变量解析、关闭后 observe() 的透传语义
  (含异常原样重新抛出)、关闭后 record_call() 零写入、关闭后重新开启恢复
  记录、4 个 admin 端点在关闭态的响应结构。

全量测试 183/183 通过(新增 10 个,含此前 173 个零回归)。

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
2026-09-14 01:10:58 +08:00
f014eb11bd HEL-543: add data hub observability side-channel (provider status, source catalog, lineage)
- New provider_call_log/provider_health tables (additive-only schema),
  wired via a fail-open observability.observe()/record_call() helper.
- Tushare pipeline keeps its existing src_calls record unchanged and now
  also feeds the unified provider_health/provider_call_log side channel.
- Eastmoney/Tencent realtime_serve.py call sites and the iFinD steward
  call site are wrapped with observability.observe() at the call site
  only; no adapter internals, routing, fallback order, or return values
  are touched.
- New read-only admin API endpoints: /admin/api/providers/status,
  /admin/api/source-catalog, /admin/api/lineage,
  /admin/api/lineage/affected.
- New static, read-only source_catalog.py and lineage.py registries
  documenting existing providers/interfaces/datasets and known
  main-site consumers (cited against backend/features/screener and
  backend/features/heaven call sites).
- provider_call_log is purged by the existing pipeline.cleanup() job
  alongside src_calls/job_runs.
- 47 new unit/integration tests covering classification, fail-open
  behavior under DB/log failures, unchanged payloads/exceptions on
  success and failure paths, and the new HTTP endpoints. Full suite:
  173 tests, all green.

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
2026-09-14 00:59:03 +08:00
8 changed files with 352 additions and 2296 deletions
+243 -1320
View File
File diff suppressed because it is too large Load Diff
+28 -89
View File
@@ -1,18 +1,16 @@
<!DOCTYPE html> <!DOCTYPE html>
<html lang="zh-CN"> <html lang="zh-CN">
<head> <head>
<meta charset="UTF-8" /> <meta charset="UTF-8" />
<meta name="viewport" content="width=device-width, initial-scale=1.0" /> <meta name="viewport" content="width=device-width, initial-scale=1" />
<title>小白复盘 · 数据中枢</title> <title>xiaobai-datahub 管理后台</title>
<link rel="stylesheet" href="/admin/styles.css" /> <link rel="stylesheet" href="/admin/styles.css" />
</head> </head>
<body> <body>
<div id="app"> <div id="app">
<section id="login-view" class="panel auth-panel">
<div id="login-view" class="auth-wrap bg-ambient" hidden>
<section class="panel auth-panel">
<h1>数据中枢</h1> <h1>数据中枢</h1>
<p class="muted">内网管理后台 · 运行总览 / 数据源配置 / 数据血缘</p> <p class="muted">内网管理后台,用于查看源状态、调度和盘后发布批次</p>
<form id="login-form"> <form id="login-form">
<label>账号 <input name="username" value="hub_admin" autocomplete="username" /></label> <label>账号 <input name="username" value="hub_admin" autocomplete="username" /></label>
<label>密码 <input name="password" type="password" autocomplete="current-password" /></label> <label>密码 <input name="password" type="password" autocomplete="current-password" /></label>
@@ -20,95 +18,36 @@
<p id="login-error" class="error" hidden></p> <p id="login-error" class="error" hidden></p>
</form> </form>
</section> </section>
</div>
<div id="change-view" class="auth-wrap bg-ambient" hidden> <section id="change-view" class="panel auth-panel" hidden>
<section class="panel auth-panel">
<h1>修改初始密码</h1> <h1>修改初始密码</h1>
<p class="muted">首次登录需先设置新密码才能进入。</p>
<form id="change-form"> <form id="change-form">
<label>当前密码 <input name="current" type="password" autocomplete="current-password" /></label> <label>当前密码 <input name="current" type="password" /></label>
<label>新密码(至少 8 位) <input name="new_password" type="password" autocomplete="new-password" /></label> <label>新密码(至少 8 位) <input name="new_password" type="password" /></label>
<button type="submit">保存并继续</button> <button type="submit">保存并继续</button>
<p id="change-error" class="error" hidden></p> <p id="change-error" class="error" hidden></p>
</form> </form>
</section> </section>
</div>
<div class="app bg-ambient" id="appRoot" hidden> <section id="shell" hidden>
<div class="bg-layer bg-gridlines z0"></div> <header class="top">
<div class="bg-layer scanlines z50"></div> <strong>xiaobai-datahub</strong>
<span id="phase" class="pill"></span>
<header class="hdr"> <span id="who" class="muted"></span>
<div class="hdr-in"> <button type="button" id="theme-btn" class="ghost">夜间</button>
<div class="radar" style="width:26px;height:26px" id="radarLogo"> <button type="button" id="logout-btn" class="ghost">退出</button>
<div class="sweep radar-sweep"></div>
<div class="ring1" style="inset:5.72px"></div>
<div class="ring2" style="inset:9.88px"></div>
<div class="center"></div>
<div class="blip pulse-dot" id="radarBlipOk" style="left:64%;top:30%;background:#34d399"></div>
<div class="blip pulse-dot" id="radarBlipBad" style="left:30%;top:62%;background:#f87171;animation-delay:.6s;display:none"></div>
</div>
<div class="brand"><span class="b1">小白复盘 <span style="color:#22d3ee">·</span> 数据中枢</span><span class="b2">DATA-HUB</span></div>
<span class="vsep"></span>
<nav class="nav" id="navEl">
<button class="navbtn" data-nav="overview"><span class="tri"></span>运行总览<span class="en">OVERVIEW</span></button>
<button class="navbtn" data-nav="sources"><span class="tri"></span>数据源配置<span class="en">SOURCES</span></button>
<button class="navbtn" data-nav="lineage"><span class="tri"></span>数据血缘<span class="en">LINEAGE</span></button>
</nav>
<span class="flex1"></span>
<span class="mdtag" id="phaseTag"></span>
<span class="livespan" id="liveSpan"><span class="livedot pulse-dot"></span>LIVE</span>
<span class="clock num" id="clock"><span id="ckD"></span><span class="csep">|</span><span class="ct"><span id="ckH"></span><span class="blink cc">:</span><span id="ckM"></span><span class="blink cc">:</span><span id="ckS"></span></span></span>
<span class="vsep"></span>
<button class="tbtn" id="opsBtn" style="font-size:10px">调度 / 发布 / 审计</button>
<button class="tbtn" id="calmBtn" style="font-size:10px">减少动态</button>
<span id="who" class="muted" style="font-size:10px"></span>
<button class="tbtn" id="logout-btn" style="font-size:10px">退出</button>
</div>
</header> </header>
<nav>
<main class="main" id="mainEl"></main> <button data-page="overview" class="active">总览</button>
<button data-page="sources">数据源</button>
<footer class="tape"> <button data-page="jobs">调度任务</button>
<span class="tape-label" id="tapeLed"><span class="led rev pulse"></span>EVENT TAPE</span> <button data-page="release">盘后发布</button>
<div class="tape-view"><div class="marquee-track" id="marqueeTrack"><span style="display:inline-flex;align-items:center" id="tapeA"></span><span style="display:inline-flex;align-items:center" id="tapeB"></span></div></div> <button data-page="datasets">数据集</button>
</footer> <button data-page="audit">审计</button>
</nav>
<main id="page"></main>
</section>
</div> </div>
<script src="/admin/app.js"></script>
<div class="drawer-mask" id="drawerMask">
<aside class="drawer">
<div class="drawer-hd">
<span class="ttl">调度 / 发布 / 审计</span>
<span class="sub" id="drawerSub"></span>
<button class="drawer-close" id="drawerClose">×</button>
</div>
<div class="drawer-tabs" id="drawerTabs">
<button class="drawer-tab" data-dtab="jobs">调度任务</button>
<button class="drawer-tab" data-dtab="release">盘后发布</button>
<button class="drawer-tab" data-dtab="audit">审计</button>
</div>
<div class="drawer-body" id="drawerBody"></div>
</aside>
</div>
<div class="modal-mask" id="modalMask">
<div class="modal-box">
<h3 id="modalTitle">危险操作确认</h3>
<p id="modalDesc"></p>
<label>管理密码 <input id="modalPassword" type="password" autocomplete="current-password" /></label>
<label id="modalConfirmWrap">请输入确认词 <span id="modalConfirmWord" class="num" style="color:#fbbf24"></span> <input id="modalConfirm" type="text" autocomplete="off" /></label>
<p class="modal-err" id="modalErr"></p>
<div class="modal-actions">
<button class="tbtn" id="modalCancel">取消</button>
<button class="tbtn danger" id="modalOk">确认执行</button>
</div>
</div>
</div>
<div class="toast" id="toastEl"></div>
</div>
<script src="/admin/app.js"></script>
</body> </body>
</html> </html>
File diff suppressed because one or more lines are too long
+5 -13
View File
@@ -31,18 +31,10 @@ class AdminAPI:
) )
is_open = bool(cal and int(cal["is_open"]) == 1) is_open = bool(cal and int(cal["is_open"]) == 1)
pubs = self.db.fetchall("SELECT * FROM publications WHERE trade_date = ?", (today,)) pubs = self.db.fetchall("SELECT * FROM publications WHERE trade_date = ?", (today,))
# HEL-529 fix 2: "待处理"只保留当前仍未恢复的最新异常。一条失败/停滞 failed = self.db.fetchall(
# 批次若已被同数据集更晚的成功批次或当日发布解决,就不再是当前故障, "SELECT * FROM batches WHERE trade_date = ? AND state IN ('failed','staged')",
# 历史记录仍完整保留在 batches/audit 明细里,不在此重复展示。 (today,),
latest_by_dataset: dict[str, dict[str, Any]] = {} )
for row in self.db.fetchall("SELECT * FROM batches WHERE trade_date = ? ORDER BY started_at, batch_id", (today,)):
latest_by_dataset[row["dataset"]] = row
published_datasets = {row["dataset"] for row in pubs if row["state"] == "published"}
anomalies = [
row
for dataset, row in latest_by_dataset.items()
if dataset not in published_datasets and row["state"] in ("failed", "staged")
]
calls = self.db.fetchall( calls = self.db.fetchall(
"SELECT * FROM src_calls ORDER BY id DESC LIMIT 20", "SELECT * FROM src_calls ORDER BY id DESC LIMIT 20",
) )
@@ -53,7 +45,7 @@ class AdminAPI:
"eod_status": self.scheduler.eod_status(today), "eod_status": self.scheduler.eod_status(today),
"revision_status": self.scheduler.revision_status(today), "revision_status": self.scheduler.revision_status(today),
"publications": pubs, "publications": pubs,
"anomalies": anomalies, "anomalies": failed,
"recent_calls": _public_calls(calls), "recent_calls": _public_calls(calls),
"source_count": len(self.db.fetchall("SELECT provider FROM src_health")), "source_count": len(self.db.fetchall("SELECT provider FROM src_health")),
} }
-17
View File
@@ -38,7 +38,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "calendar", "dataset": "calendar",
"tier": "official", "tier": "official",
"update_freq": "每日 08:45 预检",
"v1_endpoint": "/v1/calendar", "v1_endpoint": "/v1/calendar",
"primary_source": "tushare:trade_cal", "primary_source": "tushare:trade_cal",
"backup_source": None, "backup_source": None,
@@ -47,7 +46,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "stocks", "dataset": "stocks",
"tier": "official", "tier": "official",
"update_freq": "每日 20:00 / 23:10",
"v1_endpoint": "/v1/stocks", "v1_endpoint": "/v1/stocks",
"primary_source": "tushare:stock_basic", "primary_source": "tushare:stock_basic",
"backup_source": None, "backup_source": None,
@@ -56,7 +54,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "daily", "dataset": "daily",
"tier": "official", "tier": "official",
"update_freq": "盘后 15:05 · 重试至 23:30",
"v1_endpoint": "/v1/bars/daily", "v1_endpoint": "/v1/bars/daily",
"primary_source": "tushare:daily", "primary_source": "tushare:daily",
"backup_source": None, "backup_source": None,
@@ -65,7 +62,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "valuation", "dataset": "valuation",
"tier": "official", "tier": "official",
"update_freq": "盘后 15:05 · 复核 20:00",
"v1_endpoint": "/v1/valuation", "v1_endpoint": "/v1/valuation",
"primary_source": "tushare:daily_basic", "primary_source": "tushare:daily_basic",
"backup_source": None, "backup_source": None,
@@ -74,7 +70,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "moneyflow", "dataset": "moneyflow",
"tier": "official", "tier": "official",
"update_freq": "盘后 15:05 · 重试至 23:30",
"v1_endpoint": "/v1/moneyflow", "v1_endpoint": "/v1/moneyflow",
"primary_source": "tushare:moneyflow", "primary_source": "tushare:moneyflow",
"backup_source": None, "backup_source": None,
@@ -83,7 +78,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "auction", "dataset": "auction",
"tier": "official", "tier": "official",
"update_freq": "盘后 15:05",
"v1_endpoint": "/v1/auction", "v1_endpoint": "/v1/auction",
"primary_source": "tushare:stk_auction", "primary_source": "tushare:stk_auction",
"backup_source": None, "backup_source": None,
@@ -92,7 +86,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "index_daily", "dataset": "index_daily",
"tier": "official", "tier": "official",
"update_freq": "盘后 15:10 · 重试至 23:30",
"v1_endpoint": "/v1/indexes/bars", "v1_endpoint": "/v1/indexes/bars",
"primary_source": "tushare:index_daily", "primary_source": "tushare:index_daily",
"backup_source": None, "backup_source": None,
@@ -101,7 +94,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "limit_events", "dataset": "limit_events",
"tier": "official", "tier": "official",
"update_freq": "盘后 16:40 · 重试至 23:30",
"v1_endpoint": "/v1/limit-events", "v1_endpoint": "/v1/limit-events",
"primary_source": "tushare:limit_list_d", "primary_source": "tushare:limit_list_d",
"backup_source": None, "backup_source": None,
@@ -110,7 +102,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "popularity", "dataset": "popularity",
"tier": "official", "tier": "official",
"update_freq": "盘后 22:40",
"v1_endpoint": "/v1/popularity", "v1_endpoint": "/v1/popularity",
"primary_source": "tushare:ths_hot+dc_hot", "primary_source": "tushare:ths_hot+dc_hot",
"backup_source": None, "backup_source": None,
@@ -119,7 +110,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "dragon_tiger", "dataset": "dragon_tiger",
"tier": "official", "tier": "official",
"update_freq": "盘后 16:45 · 重试至 23:30",
"v1_endpoint": "/v1/dragon-tiger", "v1_endpoint": "/v1/dragon-tiger",
"primary_source": "tushare:hm_detail", "primary_source": "tushare:hm_detail",
"backup_source": None, "backup_source": None,
@@ -128,7 +118,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "sector_daily", "dataset": "sector_daily",
"tier": "official", "tier": "official",
"update_freq": "盘后 15:20 · 重试至 23:30",
"v1_endpoint": "/v1/sectors", "v1_endpoint": "/v1/sectors",
"primary_source": "tushare:ths_daily+dc_index+sw_daily", "primary_source": "tushare:ths_daily+dc_index+sw_daily",
"backup_source": None, "backup_source": None,
@@ -137,7 +126,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "quotes_latest", "dataset": "quotes_latest",
"tier": "provisional", "tier": "provisional",
"update_freq": "盘中 · 缓存 60s",
"v1_endpoint": "/v1/quotes/latest", "v1_endpoint": "/v1/quotes/latest",
"primary_source": "eastmoney:ulist/clist", "primary_source": "eastmoney:ulist/clist",
"backup_source": "tencent:qt", "backup_source": "tencent:qt",
@@ -146,7 +134,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "index_quotes", "dataset": "index_quotes",
"tier": "provisional", "tier": "provisional",
"update_freq": "盘中 · 缓存 60s",
"v1_endpoint": "/v1/indexes/quotes", "v1_endpoint": "/v1/indexes/quotes",
"primary_source": "eastmoney:ulist", "primary_source": "eastmoney:ulist",
"backup_source": "tencent:qt", "backup_source": "tencent:qt",
@@ -155,7 +142,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "sectors_quote", "dataset": "sectors_quote",
"tier": "provisional", "tier": "provisional",
"update_freq": "盘中 · 缓存 60s",
"v1_endpoint": "/v1/sectors/quote", "v1_endpoint": "/v1/sectors/quote",
"primary_source": "eastmoney:sw", "primary_source": "eastmoney:sw",
"backup_source": None, "backup_source": None,
@@ -164,7 +150,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "limit_pool", "dataset": "limit_pool",
"tier": "provisional", "tier": "provisional",
"update_freq": "盘中 · 缓存 60s",
"v1_endpoint": "/v1/limit-pool", "v1_endpoint": "/v1/limit-pool",
"primary_source": "eastmoney:zt_pool", "primary_source": "eastmoney:zt_pool",
"backup_source": None, "backup_source": None,
@@ -173,7 +158,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "intraday_points", "dataset": "intraday_points",
"tier": "provisional", "tier": "provisional",
"update_freq": "盘中 · 缓存 20s",
"v1_endpoint": "/v1/intraday/points", "v1_endpoint": "/v1/intraday/points",
"primary_source": "eastmoney:trends2", "primary_source": "eastmoney:trends2",
"backup_source": None, "backup_source": None,
@@ -182,7 +166,6 @@ DATASETS: list[dict[str, Any]] = [
{ {
"dataset": "ifind_wencai", "dataset": "ifind_wencai",
"tier": "licensed", "tier": "licensed",
"update_freq": "按需调用",
"v1_endpoint": "/v1/query (api_name=ifind_wencai)", "v1_endpoint": "/v1/query (api_name=ifind_wencai)",
"primary_source": "ifind:smart_stock_picking", "primary_source": "ifind:smart_stock_picking",
"backup_source": None, "backup_source": None,
-53
View File
@@ -99,58 +99,6 @@ STAGING_INSERT = {
**EXTENDED_STAGING_INSERT, **EXTENDED_STAGING_INSERT,
} }
# Staging-table business keys (PRIMARY KEY minus the constant batch_id),
# mirroring the PRIMARY KEY clauses declared in datahub/db.py and
# datahub/datasets_ext.py. Used only to collapse within-batch duplicates so
# a single upstream response cannot fail the whole batch on a UNIQUE
# constraint (HEL-529: 2026-09-14 dc_hot returned 4 duplicate ts_codes and
# hm_detail 43 duplicate (ts_code, hm_name) keys in one response, which has
# blocked popularity/dragon_tiger publishing every day since 09-07).
STAGING_KEY_FIELDS = {
"stocks": ("ts_code", "trade_date"),
"daily": ("ts_code", "trade_date"),
"valuation": ("ts_code", "trade_date"),
"moneyflow": ("ts_code", "trade_date"),
"auction": ("ts_code", "trade_date"),
"index_daily": ("ts_code", "trade_date"),
"limit_events": ("ts_code", "trade_date", "limit_type"),
"popularity": ("ts_code", "trade_date", "source"),
"dragon_tiger": ("ts_code", "trade_date", "hm_name"),
"sector_daily": ("ts_code", "trade_date", "family"),
}
def _dedupe_staging_rows(dataset: str, rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
"""Collapse within-batch duplicates on the staging table's business key.
Deterministic: keeps the LAST occurrence of each key (the same row the
eod copy's INSERT OR REPLACE would keep), preserves first-seen order, and
never touches rows across batches. Datasets without a declared business
key are returned unchanged.
"""
fields = STAGING_KEY_FIELDS.get(dataset)
if not fields:
return rows
seen: dict[tuple, int] = {}
out: list[dict[str, Any]] = []
dropped = 0
for row in rows:
key = tuple(row.get(field) for field in fields)
if key in seen:
out[seen[key]] = row
dropped += 1
else:
seen[key] = len(out)
out.append(row)
if dropped:
LOGGER.info(
"staging dedupe: %s collapsed %d duplicate rows within batch (kept last)",
dataset,
dropped,
extra={"hub": {"dataset": dataset, "deduped": dropped}},
)
return out
EOD_COPY = { EOD_COPY = {
"stocks": ( "stocks": (
"INSERT OR REPLACE INTO eod_stocks " "INSERT OR REPLACE INTO eod_stocks "
@@ -1772,7 +1720,6 @@ class Pipeline:
def _stage(self, dataset: str, batch_id: str, rows: list[dict[str, Any]]) -> None: def _stage(self, dataset: str, batch_id: str, rows: list[dict[str, Any]]) -> None:
sql, mapper = STAGING_INSERT[dataset] sql, mapper = STAGING_INSERT[dataset]
rows = _dedupe_staging_rows(dataset, rows)
with self.db.write() as connection: with self.db.write() as connection:
connection.execute( connection.execute(
f"DELETE FROM {DATASET_TABLES[dataset][1]} WHERE batch_id = ?", f"DELETE FROM {DATASET_TABLES[dataset][1]} WHERE batch_id = ?",
+29 -60
View File
@@ -31,21 +31,21 @@ CATALOG: list[dict[str, Any]] = [
"credential_key": "tushare_token", "credential_key": "tushare_token",
"status_source": "src_health (legacy, kept) + provider_health (unified, HEL-543)", "status_source": "src_health (legacy, kept) + provider_health (unified, HEL-543)",
"interfaces": [ "interfaces": [
{"interface": "trade_cal", "capability": "交易日历", "datasets": ["calendar"], "group": "日历 / 主档"}, {"interface": "trade_cal", "capability": "交易日历", "datasets": ["calendar"]},
{"interface": "stock_basic", "capability": "股票主档", "datasets": ["stocks"], "group": "日历 / 主档"}, {"interface": "stock_basic", "capability": "股票主档", "datasets": ["stocks"]},
{"interface": "daily", "capability": "个股日K", "datasets": ["daily"], "group": "盘后 A 批"}, {"interface": "daily", "capability": "个股日K", "datasets": ["daily"]},
{"interface": "adj_factor", "capability": "复权因子", "datasets": ["daily"], "group": "盘后 A 批"}, {"interface": "adj_factor", "capability": "复权因子", "datasets": ["daily"]},
{"interface": "daily_basic", "capability": "估值", "datasets": ["valuation"], "group": "盘后 A 批"}, {"interface": "daily_basic", "capability": "估值", "datasets": ["valuation"]},
{"interface": "index_daily", "capability": "指数日K", "datasets": ["index_daily"], "group": "指数 B 批"}, {"interface": "index_daily", "capability": "指数日K", "datasets": ["index_daily"]},
{"interface": "moneyflow", "capability": "资金流", "datasets": ["moneyflow"], "group": "盘后 A 批"}, {"interface": "moneyflow", "capability": "资金流", "datasets": ["moneyflow"]},
{"interface": "stk_auction", "capability": "集合竞价", "datasets": ["auction"], "group": "盘后 A 批"}, {"interface": "stk_auction", "capability": "集合竞价", "datasets": ["auction"]},
{"interface": "limit_list_d", "capability": "涨跌停池", "datasets": ["limit_events"], "group": "扩展软批"}, {"interface": "limit_list_d", "capability": "涨跌停池", "datasets": ["limit_events"]},
{"interface": "ths_hot", "capability": "同花顺人气榜", "datasets": ["popularity"], "group": "扩展软批"}, {"interface": "ths_hot", "capability": "同花顺人气榜", "datasets": ["popularity"]},
{"interface": "dc_hot", "capability": "东方财富人气榜", "datasets": ["popularity"], "group": "扩展软批"}, {"interface": "dc_hot", "capability": "东方财富人气榜", "datasets": ["popularity"]},
{"interface": "hm_detail", "capability": "龙虎榜游资明细", "datasets": ["dragon_tiger"], "group": "扩展软批"}, {"interface": "hm_detail", "capability": "龙虎榜游资明细", "datasets": ["dragon_tiger"]},
{"interface": "ths_daily", "capability": "同花顺概念行情", "datasets": ["sector_daily"], "group": "扩展软批"}, {"interface": "ths_daily", "capability": "同花顺概念行情", "datasets": ["sector_daily"]},
{"interface": "dc_index", "capability": "东方财富板块行情", "datasets": ["sector_daily"], "group": "扩展软批"}, {"interface": "dc_index", "capability": "东方财富板块行情", "datasets": ["sector_daily"]},
{"interface": "sw_daily", "capability": "申万行业行情", "datasets": ["sector_daily"], "group": "扩展软批"}, {"interface": "sw_daily", "capability": "申万行业行情", "datasets": ["sector_daily"]},
], ],
}, },
{ {
@@ -55,13 +55,13 @@ CATALOG: list[dict[str, Any]] = [
"credential_key": None, "credential_key": None,
"status_source": "provider_health (unified, HEL-543)", "status_source": "provider_health (unified, HEL-543)",
"interfaces": [ "interfaces": [
{"interface": "indices", "capability": "指数实时报价", "datasets": ["index_quotes"], "group": "实时快照"}, {"interface": "indices", "capability": "指数实时报价", "datasets": ["index_quotes"]},
{"interface": "market_quotes", "capability": "全市场实时快照", "datasets": ["quotes_latest"], "group": "实时快照"}, {"interface": "market_quotes", "capability": "全市场实时快照", "datasets": ["quotes_latest"]},
{"interface": "named_quotes", "capability": "指定个股实时报价", "datasets": ["quotes_latest"], "group": "实时快照"}, {"interface": "named_quotes", "capability": "指定个股实时报价", "datasets": ["quotes_latest"]},
{"interface": "sector_quote", "capability": "申万板块实时报价(单个)", "datasets": ["sectors_quote"], "group": "实时快照"}, {"interface": "sector_quote", "capability": "申万板块实时报价(单个)", "datasets": ["sectors_quote"]},
{"interface": "sector_quotes_batch", "capability": "申万板块批量报价(预热)", "datasets": ["sectors_quote"], "group": "实时快照"}, {"interface": "sector_quotes_batch", "capability": "申万板块批量报价(预热)", "datasets": ["sectors_quote"]},
{"interface": "limit_pool", "capability": "涨停/炸板池(盘中)", "datasets": ["limit_pool"], "group": "实时快照"}, {"interface": "limit_pool", "capability": "涨停/炸板池(盘中)", "datasets": ["limit_pool"]},
{"interface": "intraday", "capability": "分时走势", "datasets": ["intraday_points"], "group": "实时快照"}, {"interface": "intraday", "capability": "分时走势", "datasets": ["intraday_points"]},
], ],
}, },
{ {
@@ -71,14 +71,13 @@ CATALOG: list[dict[str, Any]] = [
"credential_key": None, "credential_key": None,
"status_source": "provider_health (unified, HEL-543)", "status_source": "provider_health (unified, HEL-543)",
"interfaces": [ "interfaces": [
{"interface": "indices", "capability": "指数实时报价(东财失败时备用)", "datasets": ["index_quotes"], "group": "实时备援"}, {"interface": "indices", "capability": "指数实时报价(东财失败时备用)", "datasets": ["index_quotes"]},
{ {
"interface": "market_quotes_fallback", "interface": "market_quotes_fallback",
"capability": "全市场快照(备用;按本地股票主档逐只请求拼接)", "capability": "全市场快照(备用;按本地股票主档逐只请求拼接)",
"datasets": ["quotes_latest"], "datasets": ["quotes_latest"],
"group": "实时备援",
}, },
{"interface": "named_quotes", "capability": "指定个股实时报价(东财失败时备用)", "datasets": ["quotes_latest"], "group": "实时备援"}, {"interface": "named_quotes", "capability": "指定个股实时报价(东财失败时备用)", "datasets": ["quotes_latest"]},
], ],
}, },
{ {
@@ -88,11 +87,11 @@ CATALOG: list[dict[str, Any]] = [
"credential_key": "ifind_refresh_token", "credential_key": "ifind_refresh_token",
"status_source": "provider_health (unified, HEL-543) + adapter.status()", "status_source": "provider_health (unified, HEL-543) + adapter.status()",
"interfaces": [ "interfaces": [
{"interface": "wencai", "capability": "问财自然语言选股", "datasets": ["ifind_wencai"], "group": "预留接口"}, {"interface": "wencai", "capability": "问财自然语言选股", "datasets": ["ifind_wencai"]},
{"interface": "snapshots", "capability": "快照", "datasets": ["ifind_snapshots"], "group": "预留接口"}, {"interface": "snapshots", "capability": "快照", "datasets": ["ifind_snapshots"]},
{"interface": "history", "capability": "历史行情", "datasets": ["ifind_history"], "group": "预留接口"}, {"interface": "history", "capability": "历史行情", "datasets": ["ifind_history"]},
{"interface": "realtime", "capability": "实时行情", "datasets": ["ifind_realtime"], "group": "预留接口"}, {"interface": "realtime", "capability": "实时行情", "datasets": ["ifind_realtime"]},
{"interface": "intraday", "capability": "分时(高频)", "datasets": ["ifind_intraday"], "group": "预留接口"}, {"interface": "intraday", "capability": "分时(高频)", "datasets": ["ifind_intraday"]},
], ],
}, },
{ {
@@ -162,35 +161,5 @@ def snapshot(db: Any, auth: Any = None) -> list[dict[str, Any]]:
except Exception: except Exception:
health_rows = [] health_rows = []
item["live_interfaces"] = health_rows item["live_interfaces"] = health_rows
# HEL-529 fix 3: 已登记 ≠ 已观测 ≠ 健康。观测记录有两种真实写法:
# eastmoney/tencent 的 observe() 直接写接口名;tushare 官方管线
# _log_call() 写的是数据集名(如 daily_basic 接口对应的数据集
# valuation)。这里按「接口名 或 该接口声明的 datasets 之一」双向
# 匹配,并为每个接口标注观测依据,禁止把"暂无观测"显示成"未配置"。
by_interface = {str(row["interface"]): row for row in health_rows}
for iface in item["interfaces"]:
row = by_interface.get(str(iface["interface"]))
basis = "interface" if row is not None else ""
if row is None:
for dataset in iface.get("datasets", []):
candidate = by_interface.get(str(dataset))
if candidate is not None:
row = candidate
basis = "dataset"
break
if row is None:
iface["observed"] = False
iface["observed_basis"] = ""
iface["observed_state"] = ""
iface["observed_at"] = ""
iface["observed_latency_ms"] = None
iface["observed_note"] = ""
else:
iface["observed"] = True
iface["observed_basis"] = basis
iface["observed_state"] = str(row["state"] or "")
iface["observed_at"] = str(row["last_ok_at"] or row["updated_at"] or "")
iface["observed_latency_ms"] = row["last_latency_ms"]
iface["observed_note"] = str(row["last_error"] or row["last_fallback_reason"] or "")
result.append(item) result.append(item)
return result return result
-213
View File
@@ -1,213 +0,0 @@
"""HEL-529 rework regressions: staging dedupe, anomaly convergence,
source-catalog observation join (dataset-name ↔ interface-name), lineage
update_freq. All read-only or within-batch fixes; none touch routing, the
8765 main site, or the 问天 frozen zone.
"""
from __future__ import annotations
import tempfile
import unittest
from pathlib import Path
from datahub.adapters.tushare import TushareAdapter
from datahub.crypto import SecretVault
from datahub.hub import Hub
from datahub.pipeline import _dedupe_staging_rows
from datahub.settings import Settings
from tests.fixtures import fake_transport
class StagingDedupeTests(unittest.TestCase):
def test_popularity_within_batch_duplicates_collapse_keep_last(self) -> None:
rows = [
{"ts_code": "600000.SH", "trade_date": "20260914", "source": "ths", "rank": 1},
{"ts_code": "000868.SZ", "trade_date": "20260914", "source": "dc", "rank": 2},
{"ts_code": "600000.SH", "trade_date": "20260914", "source": "dc", "rank": 3, "hot": 9.9},
{"ts_code": "600000.SH", "trade_date": "20260914", "source": "dc", "rank": 4, "hot": 8.8},
]
out = _dedupe_staging_rows("popularity", rows)
self.assertEqual(len(out), 3) # (600000,ths) (000868,dc) (600000,dc)
dup = [r for r in out if r["ts_code"] == "600000.SH" and r["source"] == "dc"][0]
self.assertEqual(dup["rank"], 4) # keeps LAST occurrence
self.assertEqual(out[0]["ts_code"], "600000.SH") # preserves first-seen order
def test_dragon_tiger_seat_duplicates_collapse(self) -> None:
rows = [
{"ts_code": "300010.SZ", "trade_date": "20260914", "hm_name": "T王", "buy_amount": 100},
{"ts_code": "300010.SZ", "trade_date": "20260914", "hm_name": "T王", "buy_amount": 200},
{"ts_code": "300010.SZ", "trade_date": "20260914", "hm_name": "T王", "buy_amount": 300},
]
out = _dedupe_staging_rows("dragon_tiger", rows)
self.assertEqual(len(out), 1)
self.assertEqual(out[0]["buy_amount"], 300)
def test_different_sources_are_not_duplicates(self) -> None:
rows = [
{"ts_code": "600000.SH", "trade_date": "20260914", "source": "ths"},
{"ts_code": "600000.SH", "trade_date": "20260914", "source": "dc"},
]
self.assertEqual(len(_dedupe_staging_rows("popularity", rows)), 2)
def test_unknown_dataset_passthrough(self) -> None:
rows = [{"a": 1}, {"a": 1}]
self.assertEqual(_dedupe_staging_rows("calendar", rows), rows)
class _Base(unittest.TestCase):
def setUp(self) -> None:
self.tmp = tempfile.TemporaryDirectory()
settings = Settings(
encryption_key=SecretVault.generate_key(),
api_token="z" * 32,
admin_password="StartPass1",
tushare_token="real-tushare-token-abcdef",
db_path=Path(self.tmp.name) / "hub.db",
scheduler_enabled=False,
)
self.hub = Hub(settings, adapter=TushareAdapter("real-tushare-token-abcdef", transport=fake_transport))
def tearDown(self) -> None:
self.tmp.cleanup()
class StagingDedupePublishTests(_Base):
def test_duplicate_popularity_and_dragon_tiger_now_publish(self) -> None:
"""Replays the 2026-09-14 production failure: within-response duplicate
keys used to abort the whole batch at the staging INSERT; with dedupe
the same upstream payload publishes."""
db = self.hub.db
trade_date = "20240902"
# dc_hot returns 600000.SH twice within one response; hm_detail returns
# the same (ts_code, hm_name) seat three times (mirrors live evidence).
popularity_rows = [
{"ts_code": "600000.SH", "trade_date": trade_date, "source": "ths",
"ts_name": "浦发银行", "rank": 1, "pct_change": 1.2, "current_price": 10.2,
"hot": 90.0, "concept": "银行", "data_type": "热股"},
{"ts_code": "600000.SH", "trade_date": trade_date, "source": "dc",
"ts_name": "浦发银行", "rank": 2, "pct_change": 1.2, "current_price": 10.2,
"hot": 80.0, "concept": "银行", "data_type": "A股市场"},
{"ts_code": "600000.SH", "trade_date": trade_date, "source": "dc",
"ts_name": "浦发银行", "rank": 3, "pct_change": 1.3, "current_price": 10.3,
"hot": 81.0, "concept": "银行", "data_type": "A股市场"},
]
dragon_rows = [
{"trade_date": trade_date, "ts_code": "600000.SH", "ts_name": "浦发银行",
"buy_amount": 100, "sell_amount": 200, "net_amount": -100,
"hm_name": "测试游资", "hm_orgs": "某某营业部", "tag": "超买"},
{"trade_date": trade_date, "ts_code": "600000.SH", "ts_name": "浦发银行",
"buy_amount": 300, "sell_amount": 0, "net_amount": 300,
"hm_name": "测试游资", "hm_orgs": "某某营业部", "tag": "超买"},
]
self.hub.pipeline._stage("popularity", "b-dup-pop", popularity_rows)
self.hub.pipeline._stage("dragon_tiger", "b-dup-dt", dragon_rows)
pop = db.fetchall("SELECT * FROM staging_popularity WHERE batch_id = 'b-dup-pop'")
dt = db.fetchall("SELECT * FROM staging_dragon_tiger WHERE batch_id = 'b-dup-dt'")
self.assertEqual(len(pop), 2)
self.assertEqual(len(dt), 1)
kept = dt[0]
self.assertEqual(kept["buy_amount"], 300)
class OverviewAnomalyConvergenceTests(_Base):
def _seed_batches(self, today: str) -> None:
db = self.hub.db
rows = [
# resolved history: earlier failures/stalls, later success
("x-daily-001", "daily", "staged", "empty official batch", "2026-09-14T15:05:00+08:00"),
("x-daily-002", "daily", "failed", "release group not switched", "2026-09-14T16:10:00+08:00"),
("x-daily-006", "daily", "published", "", "2026-09-14T20:00:00+08:00"),
# current unresolved faults
("x-pop-001", "popularity", "failed", "UNIQUE constraint failed: staging_popularity", "2026-09-14T22:40:00+08:00"),
("x-dt-001", "dragon_tiger", "failed", "UNIQUE constraint failed: staging_dragon_tiger", "2026-09-14T16:45:00+08:00"),
("x-dt-002", "dragon_tiger", "failed", "UNIQUE constraint failed: staging_dragon_tiger", "2026-09-14T21:48:00+08:00"),
# staged-empty later published
("x-idx-001", "index_daily", "staged", "empty official batch", "2026-09-14T15:10:00+08:00"),
("x-idx-003", "index_daily", "published", "", "2026-09-14T16:10:00+08:00"),
]
for batch_id, dataset, state, error, started in rows:
db.execute(
"INSERT INTO batches(batch_id, dataset, trade_date, state, attempt, rows_in, rows_out,"
" quality_json, started_at, finished_at, error) VALUES (?,?,?,?,?,?,?,?,?,?,?)",
(batch_id, dataset, today, state, 1, None, None, None, started, started if state != "staged" else None, error or None),
)
db.execute(
"INSERT INTO publications(dataset, trade_date, active_batch, prev_batch, state, published_at)"
" VALUES ('daily', ?, 'x-daily-006', 'x-daily-005', 'published', '2026-09-14T20:00:22+08:00')",
(today,),
)
db.execute(
"INSERT INTO publications(dataset, trade_date, active_batch, prev_batch, state, published_at)"
" VALUES ('index_daily', ?, 'x-idx-003', 'x-idx-002', 'published', '2026-09-14T16:10:12+08:00')",
(today,),
)
def test_anomalies_only_latest_unresolved(self) -> None:
overview = self.hub.admin.overview()
today = overview["trade_date"]
self._seed_batches(today)
anomalies = self.hub.admin.overview()["anomalies"]
got = sorted((a["dataset"], a["batch_id"]) for a in anomalies)
self.assertEqual(
got,
[
("dragon_tiger", "x-dt-002"), # latest failed, never published
("popularity", "x-pop-001"), # latest failed, never published
],
)
class SourceCatalogJoinTests(_Base):
def test_observation_join_matches_dataset_named_health_rows(self) -> None:
db = self.hub.db
now = "2026-09-15T08:00:00+08:00"
# tushare observability writes dataset names (real legacy behavior)
for iface, state in [("valuation", "ok"), ("popularity", "ok"), ("stocks", "ok")]:
db.execute(
"INSERT INTO provider_health(provider, interface, state, last_ok_at, last_error,"
" last_fallback_reason, consec_failures, last_latency_ms, last_data_age_seconds, updated_at)"
" VALUES ('tushare', ?, ?, ?, '', '', 0, 300, NULL, ?)",
(iface, state, now, now),
)
# realtime providers write interface names
db.execute(
"INSERT INTO provider_health(provider, interface, state, last_ok_at, last_error,"
" last_fallback_reason, consec_failures, last_latency_ms, last_data_age_seconds, updated_at)"
" VALUES ('eastmoney', 'indices', 'ok', ?, '', '', 0, 153, NULL, ?)",
(now, now),
)
items = self.hub.admin.source_catalog()["items"]
tushare = [i for i in items if i["provider"] == "tushare"][0]
by_iface = {i["interface"]: i for i in tushare["interfaces"]}
# daily_basic serves valuation → observed via dataset name
self.assertTrue(by_iface["daily_basic"]["observed"])
self.assertEqual(by_iface["daily_basic"]["observed_basis"], "dataset")
self.assertEqual(by_iface["daily_basic"]["observed_state"], "ok")
# ths_hot + dc_hot serve popularity → observed via dataset name
self.assertTrue(by_iface["ths_hot"]["observed"])
self.assertTrue(by_iface["dc_hot"]["observed"])
# stock_basic serves stocks
self.assertTrue(by_iface["stock_basic"]["observed"])
# never-observed interface stays honestly unobserved (not "unconfigured")
self.assertFalse(by_iface["trade_cal"]["observed"])
self.assertEqual(by_iface["trade_cal"]["observed_state"], "")
# interfaces carry real batch groups
groups = {i["interface"]: i["group"] for i in tushare["interfaces"]}
self.assertEqual(groups["daily"], "盘后 A 批")
self.assertEqual(groups["ths_hot"], "扩展软批")
self.assertEqual(groups["index_daily"], "指数 B 批")
eastmoney = [i for i in items if i["provider"] == "eastmoney"][0]
em = {i["interface"]: i for i in eastmoney["interfaces"]}
self.assertTrue(em["indices"]["observed"])
self.assertEqual(em["indices"]["observed_basis"], "interface")
self.assertFalse(em["market_quotes"]["observed"])
def test_lineage_update_freq_present(self) -> None:
items = self.hub.admin.lineage("20240902")["items"]
self.assertTrue(items)
for item in items:
self.assertTrue(item.get("update_freq"), f"missing update_freq for {item['dataset']}")
if __name__ == "__main__":
unittest.main()