feat(HEL-529): 按定稿100%重做三页数据中枢 + 数据修正1-5

视觉:admin/app.js 按已确认打样 hub-kimi.html 逐行重写三页 DOM 与动效
(sparkline/缓存年龄秒增/延迟变化闪烁/EVENT TAPE 预装滚动/分组接口表/
更新频率列/时钟冒号 blink/雷达 blip),CSS 补真实调用脉冲 node-ping。
数据修正:
1) pipeline._stage 批内按暂存表业务键确定性去重(保留最后一条),
   修复人气榜/龙虎榜自 09-07 起每日 UNIQUE constraint 落库失败;
2) overview.anomalies 收敛为「最新批次未成功且当日未发布」的当前异常,
   历史已恢复批次留在审计明细;
3) source_catalog 接口补真实批次分组 + 观测 join(接口名或数据集名
   双向匹配,标注 observed/observed_basis),消除「全部未配置/0/30」误报;
4) lineage 逐数据集按其服务接口过滤健康行(接口级状态),
   provisional 已配置无观测显示「已配置 · 待观测」,仅 iFinD 为未配置;
5) lineage 补 update_freq 真实频率字段。
测试:新增 tests/test_hel529_rework.py(8 项),全套 191 项通过
(1 项环境依赖失败在基线 d9358ab 上同样复现,与本改动无关)。

Co-authored-by: multica-agent <github@multica.ai>
This commit is contained in:
施工员2号
2026-09-15 08:57:29 +08:00
co-authored by multica-agent
parent d9358ab377
commit 3bafd30aad
8 changed files with 826 additions and 219 deletions
+457 -182
View File
@@ -1,12 +1,13 @@
'use strict'; 'use strict';
/* ===================================================================== /* =====================================================================
小白复盘 · 数据中枢 admin/app.js — HEL-529 第十二版 小白复盘 · 数据中枢 admin/app.js — HEL-529 返工:按定稿 100% 还原
三页真实数据中枢:运行总览 / 数据源配置 / 数据血缘。 视觉与 DOM 结构逐行照抄已确认打样 hub-kimi.htmlv12 血缘压缩增补),
视觉系统逐行照抄已通过打样 hub-kimi.html + SPEC-V11.mdv12 血缘压缩增补),
但本文件不含任何随机/模拟数据——全部状态、数字、事件均来自真实接口: 但本文件不含任何随机/模拟数据——全部状态、数字、事件均来自真实接口:
既有 /admin/api/{overview,sources,jobs,batches,datasets,audit,...} /admin/api/{overview,sources,jobs,batches,audit,...} 与 HEL-543 只读旁路
与 HEL-543 新增只读旁路 /admin/api/{providers/status,source-catalog, /admin/api/{providers/status,source-catalog,lineage,lineage/affected}。
lineage,lineage/affected}。 动效遵循 SPEC-V11marquee 46s 线性 / flowline 1.4s / radar 4.2s /
blink 1.1s steps(2) / ledPing 1.6s / sparkBlink 1.4s / 时钟冒号 blink /
延迟仅在真实值变化时闪烁(禁止随机伪造数值)/ 缓存年龄每秒真实增长。
===================================================================== */ ===================================================================== */
/* ================= 基础:DOM helper / API / CSRF ================= */ /* ================= 基础:DOM helper / API / CSRF ================= */
@@ -136,27 +137,23 @@ function kpiHtml(label, value, sub, tone, live) {
const tc = { txt: '', mint: 'glow-mint', amb: 'glow-amb', rd: 'glow-rd', cy: 'glow-cy' }[tone || 'txt']; const tc = { txt: '', mint: 'glow-mint', amb: 'glow-amb', rd: 'glow-rd', cy: 'glow-cy' }[tone || 'txt'];
return `<div class="kpi"><span class="lab">${label}</span><span class="num kv ${tc}">${value}${live ? '<span class="blink" style="color:#22d3ee;margin-left:2px">_</span>' : ''}</span>${sub ? `<span class="ks">${sub}</span>` : ''}</div>`; return `<div class="kpi"><span class="lab">${label}</span><span class="num kv ${tc}">${value}${live ? '<span class="blink" style="color:#22d3ee;margin-left:2px">_</span>' : ''}</span>${sub ? `<span class="ks">${sub}</span>` : ''}</div>`;
} }
function radarBlipsHtml(pulseOk) { /* 延迟:真实值 + 变化时才闪(打样 jit 的真实数据版:不随机抖动数字,
return pulseOk 仅当新一轮真实调用改写该值时给 .flash 青色闪烁)。 */
? `<div class="blip pulse-dot" style="left:64%;top:30%;background:#34d399"></div><div class="blip pulse-dot" style="left:30%;top:62%;background:#f87171;animation-delay:.6s"></div>`
: '';
}
function latHtml(ms, warn) { function latHtml(ms, warn) {
if (ms == null) return '<span style="color:#54637e">—</span>'; if (ms == null) return '<span style="color:#54637e">—</span>';
const color = ms > 1000 ? '#fbbf24' : (warn ? '#fbbf24' : '#d7e1f0'); const color = ms > 1000 ? '#fbbf24' : (warn ? '#fbbf24' : '#d7e1f0');
return `<span class="num" style="color:${color}">${Number(ms).toLocaleString('en-US')}ms</span>`; return `<span class="num jit" data-lat="${Number(ms)}" style="color:${color}">${Number(ms).toLocaleString('en-US')}ms</span>`;
}
function ageFromIso(iso) {
if (!iso) return null;
const t = Date.parse(iso.includes('T') ? iso : iso.replace(' ', 'T'));
if (Number.isNaN(t)) return null;
return Math.max(0, Math.floor((Date.now() - t) / 1000));
} }
/* 缓存年龄:记录真实时间戳,由 updateDynamics 每秒重算(真实增长)。 */
function ageHtml(iso) { function ageHtml(iso) {
const s = ageFromIso(iso); const ms = Date.parse(String(iso || '').includes('T') ? iso : String(iso || '').replace(' ', 'T'));
if (s == null) return '<span style="color:#54637e">—</span>'; if (!iso || Number.isNaN(ms)) return '<span style="color:#54637e">—</span>';
const txt = s < 90 ? `${s}s` : s < 3600 ? `${Math.floor(s / 60)}m` : `${Math.floor(s / 3600)}h`; return `<span class="num" data-age-epoch="${ms}" style="color:#8b9bb4">—</span>`;
return `<span class="num" style="color:${s > 90 ? '#fbbf24' : '#8b9bb4'}">${txt}</span>`; }
function fmtAgeSec(s) {
if (s == null) return '—';
// 打样口径:缓存年龄一律按秒展示、每秒增长(>30s 变琥珀)
return `${Number(s).toLocaleString('en-US')}s`;
} }
function hm(iso) { function hm(iso) {
if (!iso) return '—'; if (!iso) return '—';
@@ -167,13 +164,66 @@ function hm(iso) {
} }
const p2 = (x) => String(x).padStart(2, '0'); const p2 = (x) => String(x).padStart(2, '0');
/* ================= sparkline(真实历史:来自 provider_call_log ================= */
const SPARK_COLORS = { cy: '#22d3ee', mint: '#34d399', amb: '#fbbf24', rd: '#f87171' };
function sparkHtml(jid, tone, w, h) {
return `<svg width="${w}" height="${h}" class="spark" data-jid="${esc(jid)}" data-w="${w}" data-h="${h}" style="display:block;opacity:.9"><polyline points="" fill="none" stroke="${SPARK_COLORS[tone]}" stroke-width="1.2" stroke-opacity=".8"/><circle class="spark-end" r="2" fill="${SPARK_COLORS[tone]}" cx="0" cy="0"/></svg>`;
}
function sparkSeries(jid) {
const hist = LatHist.get(jid) || [];
if (hist.length < 2) return null;
const max = Math.max(...hist, 1);
return hist.map((v) => v / max);
}
function sparkPaint(svg) {
const jid = svg.dataset.jid, w = +svg.dataset.w, h = +svg.dataset.h;
const data = sparkSeries(jid);
const poly = svg.querySelector('polyline');
const dot = svg.querySelector('circle');
if (!data) {
const y = (h / 2).toFixed(2);
poly.setAttribute('points', `0,${y} ${w},${y}`);
poly.setAttribute('stroke-opacity', '.25');
dot.style.display = 'none';
return;
}
poly.setAttribute('stroke-opacity', '.8');
dot.style.display = '';
const pts = data.map((v, i) => `${((i / (data.length - 1)) * w).toFixed(2)},${(h - v * (h - 3) - 1).toFixed(2)}`).join(' ');
poly.setAttribute('points', pts);
dot.setAttribute('cx', (w - 1).toFixed(2));
dot.setAttribute('cy', (h - data[data.length - 1] * (h - 3) - 1).toFixed(2));
}
/* 每个“观测对象”(接口或数据集)的真实延迟历史(旧→新),来自
provider_call_log 最近 200 条。 */
const LatHist = new Map();
const CALL_MATCHERS = {};
function buildLatHist() {
const calls = (RAW.providers && RAW.providers.recent_calls) || [];
const buckets = new Map();
[...calls].reverse().forEach((c) => {
if (c.latency_ms == null) return;
const key = `${c.provider}:${c.interface}`;
(buckets.get(key) || buckets.set(key, []).get(key)).push(c.latency_ms);
});
LatHist.clear();
buckets.forEach((arr, key) => LatHist.set(key, arr.slice(-24)));
// 数据集维度的历史:tushare 观测行写的是数据集名,直接按数据集聚合
const dsBuckets = new Map();
[...calls].reverse().forEach((c) => {
if (c.provider !== 'tushare' || c.latency_ms == null) return;
(dsBuckets.get(c.interface) || dsBuckets.set(c.interface, []).get(c.interface)).push(c.latency_ms);
});
dsBuckets.forEach((arr, ds) => LatHist.set(`ds:${ds}`, arr.slice(-24)));
}
/* ================= 数据层:拉取真实接口并汇总为一份快照 ================= */ /* ================= 数据层:拉取真实接口并汇总为一份快照 ================= */
const RAW = { const RAW = {
overview: null, sources: null, providers: null, catalog: null, overview: null, sources: null, providers: null, catalog: null,
lineage: null, jobs: null, batches: null, datasets: null, audit: null, lineage: null, jobs: null, batches: null, audit: null,
}; };
async function fetchAll() { async function fetchAll() {
const [overview, sources, providers, catalog, lineage, jobs, batches, datasets, audit] = await Promise.all([ const [overview, sources, providers, catalog, lineage, jobs, batches, audit] = await Promise.all([
api('/admin/api/overview'), api('/admin/api/overview'),
api('/admin/api/sources'), api('/admin/api/sources'),
api('/admin/api/providers/status?limit=200'), api('/admin/api/providers/status?limit=200'),
@@ -181,31 +231,14 @@ async function fetchAll() {
api('/admin/api/lineage'), api('/admin/api/lineage'),
api('/admin/api/jobs'), api('/admin/api/jobs'),
api(`/admin/api/batches?date=${today()}`), api(`/admin/api/batches?date=${today()}`),
api(`/admin/api/datasets?date=${today()}`),
api('/admin/api/audit'), api('/admin/api/audit'),
]); ]);
RAW.overview = overview; RAW.sources = sources; RAW.providers = providers; RAW.catalog = catalog; RAW.overview = overview; RAW.sources = sources; RAW.providers = providers; RAW.catalog = catalog;
RAW.lineage = lineage; RAW.jobs = jobs; RAW.batches = batches; RAW.datasets = datasets; RAW.audit = audit; RAW.lineage = lineage; RAW.jobs = jobs; RAW.batches = batches; RAW.audit = audit;
buildLatHist();
} }
function today() { return RAW.overview ? RAW.overview.trade_date : ''; } function today() { return RAW.overview ? RAW.overview.trade_date : ''; }
/* 按 provider(+interface) 索引 provider_health,供各页复用 */
function healthIndex() {
const idx = {};
const rows = (RAW.providers && RAW.providers.health) || [];
for (const row of rows) {
idx[row.provider] = idx[row.provider] || [];
idx[row.provider].push(row);
}
return idx;
}
function callsIndex() {
const idx = {};
const rows = (RAW.providers && RAW.providers.recent_calls) || [];
for (const row of rows) { (idx[row.provider] = idx[row.provider] || []).push(row); }
return idx;
}
const OFFICIAL_DATASET_IDS = ['daily', 'valuation', 'moneyflow', 'auction', 'index_daily', 'limit_events', 'popularity', 'dragon_tiger', 'sector_daily', 'stocks']; const OFFICIAL_DATASET_IDS = ['daily', 'valuation', 'moneyflow', 'auction', 'index_daily', 'limit_events', 'popularity', 'dragon_tiger', 'sector_daily', 'stocks'];
const DATASET_NAME = { const DATASET_NAME = {
daily: '日K', valuation: '估值', moneyflow: '资金流', auction: '竞价', index_daily: '指数日K', daily: '日K', valuation: '估值', moneyflow: '资金流', auction: '竞价', index_daily: '指数日K',
@@ -213,7 +246,7 @@ const DATASET_NAME = {
}; };
/* ===================================================================== /* =====================================================================
总览页 总览页(结构照抄 hub-kimi.html,数据全部真实)
===================================================================== */ ===================================================================== */
function buildOverviewVM() { function buildOverviewVM() {
const ov = RAW.overview || {}; const ov = RAW.overview || {};
@@ -223,7 +256,6 @@ function buildOverviewVM() {
(ov.publications || []).forEach((p) => { pubsByDs[p.dataset] = p; }); (ov.publications || []).forEach((p) => { pubsByDs[p.dataset] = p; });
const batchByDs = {}; const batchByDs = {};
((RAW.batches && RAW.batches.batches) || []).forEach((b) => { if (!batchByDs[b.dataset] || b.started_at > batchByDs[b.dataset].started_at) batchByDs[b.dataset] = b; }); ((RAW.batches && RAW.batches.batches) || []).forEach((b) => { if (!batchByDs[b.dataset] || b.started_at > batchByDs[b.dataset].started_at) batchByDs[b.dataset] = b; });
const hidx = healthIndex();
const totalOfficial = OFFICIAL_DATASET_IDS.length; const totalOfficial = OFFICIAL_DATASET_IDS.length;
const publishedCount = OFFICIAL_DATASET_IDS.filter((id) => pubsByDs[id] && pubsByDs[id].state === 'published').length; const publishedCount = OFFICIAL_DATASET_IDS.filter((id) => pubsByDs[id] && pubsByDs[id].state === 'published').length;
const anomalyCount = (ov.anomalies || []).length; const anomalyCount = (ov.anomalies || []).length;
@@ -231,7 +263,7 @@ function buildOverviewVM() {
const successRate = recentCalls.length ? (recentCalls.filter((c) => c.ok).length / recentCalls.length * 100) : null; const successRate = recentCalls.length ? (recentCalls.filter((c) => c.ok).length / recentCalls.length * 100) : null;
// 主站回退页面:主源异常且备源健康 → 判定为“回退中”的 provisional 数据集 // 主站回退页面:主源异常且备源健康 → 判定为“回退中”的 provisional 数据集
let fallbackPages = 0; const fallbackPages = [];
(RAW.lineage && RAW.lineage.items || []).forEach((it) => { (RAW.lineage && RAW.lineage.items || []).forEach((it) => {
if (it.tier !== 'provisional' || !it.backup_source) return; if (it.tier !== 'provisional' || !it.backup_source) return;
const rows = it.live_provider_health || []; const rows = it.live_provider_health || [];
@@ -239,7 +271,9 @@ function buildOverviewVM() {
const backupProv = (it.providers || [])[1]; const backupProv = (it.providers || [])[1];
const primaryState = worstBucket(rows.filter((r) => r.provider === primaryProv).map((r) => bucketOf(r.state))); const primaryState = worstBucket(rows.filter((r) => r.provider === primaryProv).map((r) => bucketOf(r.state)));
const backupState = worstBucket(rows.filter((r) => r.provider === backupProv).map((r) => bucketOf(r.state))); const backupState = worstBucket(rows.filter((r) => r.provider === backupProv).map((r) => bucketOf(r.state)));
if ((primaryState === 'fail' || primaryState === 'degrade' || primaryState === 'slow') && backupState === 'ok') fallbackPages += 1; if ((primaryState === 'fail' || primaryState === 'degrade' || primaryState === 'slow') && backupState === 'ok') {
(it.known_consumers || []).forEach((c) => fallbackPages.push(c));
}
}); });
const datasetCards = OFFICIAL_DATASET_IDS.map((id) => { const datasetCards = OFFICIAL_DATASET_IDS.map((id) => {
@@ -257,37 +291,31 @@ function buildOverviewVM() {
return { id, name: DATASET_NAME[id] || id, bucket, time: note && !batch?.error ? note : '', rows, err: batch && batch.error ? batch.error : '' }; return { id, name: DATASET_NAME[id] || id, bucket, time: note && !batch?.error ? note : '', rows, err: batch && batch.error ? batch.error : '' };
}); });
// 实时观察层:eastmoney / tencent 的真实 provider_health 行 // 实时观察层:eastmoney / tencent 目录接口的真实观测(observed=true
const catalogByProv = {}; const catalogByProv = {};
(RAW.catalog && RAW.catalog.items || []).forEach((it) => { catalogByProv[it.provider] = it; }); (RAW.catalog && RAW.catalog.items || []).forEach((it) => { catalogByProv[it.provider] = it; });
const observers = []; const observers = [];
['eastmoney', 'tencent'].forEach((prov) => { ['eastmoney', 'tencent'].forEach((prov) => {
const rows = hidx[prov] || [];
const cat = catalogByProv[prov]; const cat = catalogByProv[prov];
rows.forEach((row) => { if (!cat) return;
const iface = (cat && cat.interfaces.find((i) => i.interface === row.interface)) || {}; cat.interfaces.forEach((iface) => {
if (!iface.observed) return;
observers.push({ observers.push({
name: iface.capability || row.interface, jid: `${prov}:${iface.interface}`,
provider: prov, name: iface.capability, provider: prov, interface: iface.interface,
interface: row.interface, lastOk: iface.observed_at, lat: iface.observed_latency_ms,
lastOk: row.last_ok_at, bucket: bucketOf(iface.observed_state),
lat: row.last_latency_ms, note: iface.observed_note || '',
bucket: bucketOf(row.state),
note: row.last_error || row.last_fallback_reason || '',
}); });
}); });
}); });
// 主站影响:lineage() 条目按 known_consumers 展开 // 主站影响:lineage() 条目按 known_consumers 展开(接口级状态)
const siteImpact = []; const siteImpact = [];
(RAW.lineage && RAW.lineage.items || []).forEach((it) => { (RAW.lineage && RAW.lineage.items || []).forEach((it) => {
const rows = it.live_provider_health || []; const bucket = lineageItemBucket(it);
const bucket = it.tier === 'official'
? (it.publication ? (it.publication.state === 'published' ? 'ok' : (it.publication.state === 'degraded' ? 'degrade' : 'fail')) : (rows.length ? worstBucket(rows.map((r) => bucketOf(r.state))) : 'plan'))
: (rows.length ? worstBucket(rows.map((r) => bucketOf(r.state))) : 'off');
const note = rows.find((r) => r.last_error)?.last_error || rows.find((r) => r.last_fallback_reason)?.last_fallback_reason || '';
(it.known_consumers || []).forEach((c) => { (it.known_consumers || []).forEach((c) => {
siteImpact.push({ page: c, dataset: it.dataset, bucket, note }); siteImpact.push({ page: c, short: pageNameOf(c), dataset: it.dataset, bucket, note: lineageItemNote(it) });
}); });
}); });
const okPages = siteImpact.filter((s) => s.bucket === 'ok').length; const okPages = siteImpact.filter((s) => s.bucket === 'ok').length;
@@ -304,14 +332,16 @@ function buildOverviewVM() {
return { t: (j.at.split('-')[0] || j.at).slice(0, 5), label: j.title, state: st, jobId: j.id }; return { t: (j.at.split('-')[0] || j.at).slice(0, 5), label: j.title, state: st, jobId: j.id };
}); });
// 待处理异常:批次失败/停滞 + 主源连续失败 + 未配置授权源 // 待处理异常:overview.anomalies 已在后端收敛为“仍未恢复的最新异常”
const incidents = []; const incidents = [];
(ov.anomalies || []).forEach((b) => { (ov.anomalies || []).forEach((b) => {
incidents.push({ sev: 'fail', tag: b.state, title: `${DATASET_NAME[b.dataset] || b.dataset} ${b.state === 'failed' ? '失败' : '停滞'}`, meta: `${b.error ? esc(b.error) : '无错误详情'} · 批次 ${b.batch_id}`, dataset: b.dataset }); incidents.push({ sev: 'fail', tag: b.state, title: `${DATASET_NAME[b.dataset] || b.dataset} ${b.state === 'failed' ? '失败' : '停滞'}`, meta: `${b.error ? esc(b.error) : '无错误详情'} · 批次 ${b.batch_id}`, dataset: b.dataset });
}); });
if (eod.state === 'waiting_upstream' || eod.state === 'cutoff_failed') { if (eod.state === 'waiting_upstream' || eod.state === 'cutoff_failed') {
incidents.push({ sev: eod.state === 'cutoff_failed' ? 'fail' : 'slow', tag: eod.state, title: eod.state === 'cutoff_failed' ? '盘后批已截止失败' : '等待上游重试', meta: `已试 ${eod.attempts || 0}${(eod.missing_datasets || []).length ? ' · 缺 ' + eod.missing_datasets.join(',') : ''}`, dataset: (eod.missing_datasets || [])[0] || '' }); incidents.push({ sev: eod.state === 'cutoff_failed' ? 'fail' : 'slow', tag: eod.state, title: eod.state === 'cutoff_failed' ? '盘后批已截止失败' : '自动重试窗口', meta: `15:1523:30 · 每 30 分钟 · 成功即停 · 已试 ${eod.attempts || 0}${(eod.missing_datasets || []).length ? ' · 缺 ' + eod.missing_datasets.join(',') : ''}`, dataset: (eod.missing_datasets || [])[0] || '', cd: eod.next_retry_at });
} }
const hidx = {};
((RAW.providers && RAW.providers.health) || []).forEach((r) => { (hidx[r.provider] = hidx[r.provider] || []).push(r); });
Object.keys(hidx).forEach((prov) => { Object.keys(hidx).forEach((prov) => {
const bad = (hidx[prov] || []).filter((r) => (r.consec_failures || 0) > 0 && bucketOf(r.state) === 'fail'); const bad = (hidx[prov] || []).filter((r) => (r.consec_failures || 0) > 0 && bucketOf(r.state) === 'fail');
if (bad.length) incidents.push({ sev: 'fail', tag: prov, title: `${prov} 部分接口连续失败`, meta: bad.map((r) => `${r.interface}×${r.consec_failures}`).join(' · '), provider: prov }); if (bad.length) incidents.push({ sev: 'fail', tag: prov, title: `${prov} 部分接口连续失败`, meta: bad.map((r) => `${r.interface}×${r.consec_failures}`).join(' · '), provider: prov });
@@ -325,6 +355,7 @@ function buildOverviewVM() {
trade_date: ov.trade_date, session_phase: ov.session_phase, eod, rev, trade_date: ov.trade_date, session_phase: ov.session_phase, eod, rev,
publishedCount, totalOfficial, anomalyCount, successRate, publishedCount, totalOfficial, anomalyCount, successRate,
fallbackPages, datasetCards, observers, siteImpact, okPages, timeline, incidents, fallbackPages, datasetCards, observers, siteImpact, okPages, timeline, incidents,
unpublished: OFFICIAL_DATASET_IDS.filter((id) => !pubsByDs[id] || pubsByDs[id].state !== 'published'),
}; };
} }
@@ -335,7 +366,7 @@ function heroHtml(vm) {
if (vm.rev.state === 'waiting_review') tags.push(tagHtml('cy', '复核 · 待发布')); if (vm.rev.state === 'waiting_review') tags.push(tagHtml('cy', '复核 · 待发布'));
const subParts = []; const subParts = [];
if (vm.eod.attempts) subParts.push(`已试 <span style="color:#fbbf24">${vm.eod.attempts}</span> 次`); if (vm.eod.attempts) subParts.push(`已试 <span style="color:#fbbf24">${vm.eod.attempts}</span> 次`);
if (vm.eod.next_retry_at) subParts.push(`下次重试 <span style="color:#fbbf24">${hm(vm.eod.next_retry_at)}</span>`); if (vm.eod.next_retry_at) subParts.push(`下次重试 <span class="num" data-cd-epoch="${Date.parse(String(vm.eod.next_retry_at).replace(' ', 'T'))}" data-cdfmt="mmss" style="color:#fbbf24">--:--</span>`);
if ((vm.eod.missing_datasets || []).length) subParts.push(`缺 <span style="color:#f87171">${esc(vm.eod.missing_datasets.join(','))}</span>`); if ((vm.eod.missing_datasets || []).length) subParts.push(`缺 <span style="color:#f87171">${esc(vm.eod.missing_datasets.join(','))}</span>`);
if (vm.rev.window) subParts.push(`复核窗口 <span style="color:#22d3ee">${esc(vm.rev.window)}</span>`); if (vm.rev.window) subParts.push(`复核窗口 <span style="color:#22d3ee">${esc(vm.rev.window)}</span>`);
return `<div class="panel hero">` + return `<div class="panel hero">` +
@@ -345,23 +376,25 @@ function heroHtml(vm) {
(subParts.length ? `<div class="hero-sub">${subParts.join('<span class="sep">|</span>')}</div>` : '') + (subParts.length ? `<div class="hero-sub">${subParts.join('<span class="sep">|</span>')}</div>` : '') +
`</div><span class="flex1"></span>` + `</div><span class="flex1"></span>` +
`<div class="kpis">` + `<div class="kpis">` +
kpiHtml('已发布数据集', `<span class="glow-mint">${vm.publishedCount}</span><span style="color:#54637e;font-size:18px">/${vm.totalOfficial}</span>`, null) + kpiHtml('已发布数据集', `<span class="glow-mint">${vm.publishedCount}</span><span style="color:#54637e;font-size:18px">/${vm.totalOfficial}</span>`, vm.unpublished.length ? `${vm.unpublished.map((d) => DATASET_NAME[d] || d).join(' · ')} 未发布` : '全部已发布') +
kpiHtml('今日调用成功率', vm.successRate == null ? '<span style="color:#54637e">—</span>' : `<span>${vm.successRate.toFixed(1)}%</span>`, null, 'mint', vm.successRate != null) + kpiHtml('今日调用成功率', vm.successRate == null ? '<span style="color:#54637e">—</span>' : `<span>${vm.successRate.toFixed(1)}%</span>`, null, 'mint', vm.successRate != null) +
kpiHtml('主站回退页面', `${vm.fallbackPages}`, null, vm.fallbackPages ? 'amb' : 'txt') + kpiHtml('主站回退页面', `${vm.fallbackPages.length}`, vm.fallbackPages.length ? esc(vm.fallbackPages[0]) : '实时层主源正常', vm.fallbackPages.length ? 'amb' : 'txt') +
kpiHtml('EOD 尝试', `${vm.eod.attempts || 0}`, '上限 5 · 23:30 止', vm.eod.attempts ? 'rd' : 'txt') + kpiHtml('EOD 尝试', `${vm.eod.attempts || 0}`, '上限 5 · 23:30 止', vm.eod.attempts ? 'rd' : 'txt') +
`</div></div>`; `</div></div>`;
} }
function datasetStripHtml(vm) { function datasetStripHtml(vm) {
return `<div class="ds-grid">${vm.datasetCards.map((d) => { return `<div class="ds-grid">${vm.datasetCards.map((d) => {
const ring = d.bucket === 'fail' ? 'fail' : d.bucket === 'degrade' ? 'slow' : d.bucket === 'review' ? 'review' : ''; const ring = d.bucket === 'fail' ? 'fail' : d.bucket === 'degrade' ? 'slow' : d.bucket === 'review' ? 'review' : '';
const sparkTone = d.bucket === 'fail' ? 'rd' : d.bucket === 'degrade' ? 'amb' : d.bucket === 'review' ? 'cy' : 'mint';
let note = ''; let note = '';
if (d.time) note += `<span class="num" style="color:#8b9bb4">${esc(d.time)} · </span>`; if (d.time) note += `<span class="num" style="color:#8b9bb4">${esc(d.time)} · </span>`;
if (d.rows) note += `<span class="num">${Number(d.rows).toLocaleString('en-US')} 行</span>`; if (d.rows) note += `<span class="num">${Number(d.rows).toLocaleString('en-US')} 行</span>`;
if (d.err) note += `<span style="color:#f87171">${esc(d.err)}</span>`; if (d.err) { const brief = String(d.err).length > 22 ? String(d.err).slice(0, 22) + '…' : d.err; note += `<span style="color:#f87171" title="${esc(d.err)}">${esc(brief)}</span>`; }
if (!d.time && !d.rows && !d.err) note = '<span style="color:#54637e">今日尚未发布</span>'; if (!d.time && !d.rows && !d.err) note = '<span style="color:#54637e">今日尚未发布</span>';
return `<div class="panel dsc ${ring}" data-dataset="${d.id}">` + return `<div class="panel dsc ${ring}" data-dataset="${d.id}">` +
`<div class="dsc-top"><span class="dsc-name">${esc(d.name)}</span>${pillHtml(d.bucket)}</div>` + `<div class="dsc-top"><span class="dsc-name">${esc(d.name)}</span>${pillHtml(d.bucket)}</div>` +
`<div class="dsc-note">${note}</div>` + `<div class="dsc-note">${note}</div>` +
sparkHtml(`ds:${d.id}`, sparkTone, 96, 16) +
`</div>`; `</div>`;
}).join('')}</div>`; }).join('')}</div>`;
} }
@@ -373,20 +406,30 @@ function observatoryHtml(vm) {
`<td style="color:#8b9bb4">${o.provider === 'eastmoney' ? '东方财富' : '腾讯行情'}</td>` + `<td style="color:#8b9bb4">${o.provider === 'eastmoney' ? '东方财富' : '腾讯行情'}</td>` +
`<td>${ageHtml(o.lastOk)}</td>` + `<td>${ageHtml(o.lastOk)}</td>` +
`<td>${latHtml(o.lat, o.bucket !== 'ok')}</td>` + `<td>${latHtml(o.lat, o.bucket !== 'ok')}</td>` +
`<td>${sparkHtml(o.jid, o.bucket === 'slow' ? 'amb' : 'mint', 64, 16)}</td>` +
`<td><span style="display:flex;align-items:center;gap:8px">${pillHtml(o.bucket)}${o.note ? `<span class="obs-note">${esc(o.note)}</span>` : ''}</span></td>` + `<td><span style="display:flex;align-items:center;gap:8px">${pillHtml(o.bucket)}${o.note ? `<span class="obs-note">${esc(o.note)}</span>` : ''}</span></td>` +
`</tr>`).join('') : `<tr><td colspan="5" class="empty-hint">尚无真实观测记录(等待首次探测/调用)</td></tr>`; `</tr>`).join('') : `<tr><td colspan="6" class="empty-hint">尚无真实观测记录(等待首次探测/调用)</td></tr>`;
const body = `<table class="dtable"><thead><tr><th>观察项</th><th>当前来源</th><th>最近成功</th><th>延迟</th><th>状态 · 备注</th></tr></thead><tbody>${rows}</tbody></table>`; const body = `<table class="dtable"><thead><tr><th>观察项</th><th>当前来源</th><th>缓存年龄</th><th>延迟</th><th>趋势</th><th>状态 · 备注</th></tr></thead><tbody>${rows}</tbody></table>`;
return panelHtml('实时观察层', right, body, false); return panelHtml('实时观察层', right, body, false);
} }
function siteImpactHtml(vm) { function siteImpactHtml(vm) {
const bad = vm.siteImpact.filter((s) => s.bucket !== 'ok');
const right = `<span style="font-size:10px"><span style="color:#34d399">${vm.okPages} 项正常</span><span style="color:#54637e"> · </span><span style="color:#fbbf24">${vm.siteImpact.length - vm.okPages} 项待关注</span></span>`; const right = `<span style="font-size:10px"><span style="color:#34d399">${vm.okPages} 项正常</span><span style="color:#54637e"> · </span><span style="color:#fbbf24">${vm.siteImpact.length - vm.okPages} 项待关注</span></span>`;
const rows = vm.siteImpact.length ? vm.siteImpact.map((s) => `<tr>` + const rows = bad.length ? bad.slice(0, 5).map((s) => `<tr>` +
`<td style="color:#e8f1ff;width:36%">${esc(s.page)}</td>` + `<td style="color:#e8f1ff;width:40%">${esc(s.page)}</td>` +
`<td style="width:18%">${pillHtml(s.bucket)}</td>` + `<td style="width:20%">${pillHtml(s.bucket)}</td>` +
`<td style="font-size:11px;color:#54637e">${s.note ? esc(s.note) : '—'}</td>` + `<td style="font-size:11px;color:#54637e">${s.note ? esc(s.note) : '—'}</td>` +
`<td style="text-align:right"><button class="tbtn mini" data-goto-lineage="${esc(s.dataset)}">查看血缘</button></td>` + `<td style="text-align:right"><button class="tbtn mini" data-goto-lineage="${esc(s.dataset)}">查看血缘</button></td>` +
`</tr>`).join('') : `<tr><td colspan="4" class="empty-hint">尚无血缘记录</td></tr>`; `</tr>`).join('') : `<tr><td style="color:#e8f1ff;width:40%">全部页面正常</td><td>${pillHtml('ok')}</td><td style="font-size:11px;color:#54637e">—</td><td></td></tr>`;
return panelHtml('主站影响', right, `<table class="dtable"><tbody>${rows}</tbody></table>`, false); const footPages = vm.siteImpact.map((s) => ({ name: s.short, bucket: s.bucket }));
const byPage = new Map();
footPages.forEach((p) => { const cur = byPage.get(p.name); if (!cur || SEV[p.bucket] > SEV[cur]) byPage.set(p.name, p.bucket); });
const foot = `<div class="site-foot"><span style="color:#34d399;margin-right:4px">正常 ${vm.okPages}</span>` +
[...byPage.entries()].map(([name, bucket]) => bucket === 'ok'
? `<span style="display:inline-flex;align-items:center;gap:6px"><span>${esc(name)}</span><span style="color:#243152">·</span></span>`
: `<span style="display:inline-flex;align-items:center;gap:6px;color:${HEALTH[bucket].color}">${esc(name)}(${HEALTH[bucket].label})</span><span style="color:#243152">·</span></span>`).join('') +
`</div>`;
return panelHtml('主站影响', right, `<table class="dtable"><tbody>${rows}</tbody></table>${foot}`, false);
} }
function timelineHtml(vm) { function timelineHtml(vm) {
const nodes = vm.timeline; const nodes = vm.timeline;
@@ -422,17 +465,21 @@ function timelineHtml(vm) {
s += `<g><line x1="${nowX}" y1="${y - 44}" x2="${nowX}" y2="${y + 44}" stroke="#22d3ee" stroke-width="1" stroke-dasharray="2 3" opacity=".7"/>` + s += `<g><line x1="${nowX}" y1="${y - 44}" x2="${nowX}" y2="${y + 44}" stroke="#22d3ee" stroke-width="1" stroke-dasharray="2 3" opacity=".7"/>` +
`<circle cx="${nowX}" cy="${y}" r="3.6" fill="#22d3ee" class="now-dot now-r"/>` + `<circle cx="${nowX}" cy="${y}" r="3.6" fill="#22d3ee" class="now-dot now-r"/>` +
`<text x="${nowX}" y="${y - 50}" text-anchor="middle" font-size="9" fill="#22d3ee" font-weight="600" class="num">NOW</text></g></svg>`; `<text x="${nowX}" y="${y - 50}" text-anchor="middle" font-size="9" fill="#22d3ee" font-weight="600" class="num">NOW</text></g></svg>`;
const right = `<span style="font-size:10px;color:#54637e">现在 <span class="num" style="color:#22d3ee">${p2(now.getHours())}:${p2(now.getMinutes())}</span> · <span style="color:#8b9bb4">SCHEDULER.PY</span></span>`; const range = nodes.length ? `${nodes[0].t}${nodes[nodes.length - 1].t}` : '';
const right = `<span style="font-size:10px;color:#54637e">现在 <span class="num" style="color:#22d3ee">${p2(now.getHours())}:${p2(now.getMinutes())}</span> · ${range} · <span style="color:#8b9bb4">SCHEDULER.PY</span></span>`;
return panelHtml('今晚时间线', right, s, false); return panelHtml('今晚时间线', right, s, false);
} }
function incidentsHtml(vm) { function incidentsHtml(vm) {
if (!vm.incidents.length) return panelHtml('待处理异常', tagHtml('mint', '0 OPEN'), '<div class="empty-hint">当前无待处理异常</div>', true); if (!vm.incidents.length) return panelHtml('待处理异常', tagHtml('mint', '0 OPEN'), '<div class="empty-hint">当前无待处理异常</div>', true);
const cards = vm.incidents.map((it) => { const cards = vm.incidents.map((it) => {
const tag = it.sev === 'fail' ? tagHtml('rd', esc(it.tag)) : it.sev === 'slow' ? tagHtml('amb', esc(it.tag)) : tagHtml('dashed', esc(it.tag)); const tag = it.sev === 'fail' ? tagHtml('rd', esc(it.tag)) : it.sev === 'slow' ? tagHtml('amb', esc(it.tag)) : tagHtml('dashed', esc(it.tag));
const meta = it.cd
? `${it.meta}<span class="num" style="color:#fbbf24"> · 下次重试 <span data-cd-epoch="${Date.parse(String(it.cd).replace(' ', 'T'))}" data-cdfmt="ms">--m--s</span></span>`
: it.meta;
return `<div class="inc ${it.sev}">${tag}` + return `<div class="inc ${it.sev}">${tag}` +
`<div style="flex:1;min-width:0">` + `<div style="flex:1;min-width:0">` +
`<div class="inc-title"${it.sev === 'off' ? ' style="color:#8b9bb4"' : ''}>${esc(it.title)}</div>` + `<div class="inc-title"${it.sev === 'off' ? ' style="color:#8b9bb4"' : ''}>${esc(it.title)}</div>` +
`<div class="inc-meta">${esc(it.meta)}</div>` + `<div class="inc-meta">${meta}</div>` +
`</div>` + `</div>` +
`<button class="tbtn" style="flex:none" data-incident-action data-dataset="${esc(it.dataset || '')}" data-provider="${esc(it.provider || '')}">去处理 →</button>` + `<button class="tbtn" style="flex:none" data-incident-action data-dataset="${esc(it.dataset || '')}" data-provider="${esc(it.provider || '')}">去处理 →</button>` +
`</div>`; `</div>`;
@@ -447,44 +494,62 @@ function overviewHtml() {
} }
/* ===================================================================== /* =====================================================================
数据源配置页 数据源配置页(分组接口表照抄打样 mx-grid 结构,数据全部真实)
===================================================================== */ ===================================================================== */
const srcUiState = {}; const srcUiState = {};
const PROV_ROLE_TEXT = {
tushare: 'official · 盘后正式数据唯一源',
eastmoney: 'free · 盘中观察主源',
tencent: 'free · 盘中观察备源',
ifind: 'licensed · 图表 / 问财 / 实时快照',
};
function buildSourcesVM() { function buildSourcesVM() {
const catalog = (RAW.catalog && RAW.catalog.items) || []; const catalog = (RAW.catalog && RAW.catalog.items) || [];
const legacy = {}; const legacy = {};
(RAW.sources && RAW.sources.items || []).forEach((it) => { legacy[it.provider] = it; }); (RAW.sources && RAW.sources.items || []).forEach((it) => { legacy[it.provider] = it; });
const hidx = healthIndex(); const cidx = {};
const cidx = callsIndex(); ((RAW.providers && RAW.providers.recent_calls) || []).forEach((c) => { if (!cidx[c.provider]) cidx[c.provider] = c; });
const active = catalog.filter((c) => c.role !== 'reserved'); const active = catalog.filter((c) => c.role !== 'reserved');
const reserved = catalog.filter((c) => c.role === 'reserved'); const reserved = catalog.filter((c) => c.role === 'reserved');
const totalIfaces = active.reduce((n, c) => n + c.interfaces.length, 0); const totalIfaces = active.reduce((n, c) => n + c.interfaces.length, 0);
const usedIfaces = active.reduce((n, c) => n + c.interfaces.filter((i) => (hidx[c.provider] || []).some((h) => h.interface === i.interface)).length, 0); // “在用”= 已登记且已观察到调用(接口名或数据集名匹配,后端 join 已标注)
const usedIfaces = active.reduce((n, c) => n + c.interfaces.filter((i) => i.observed).length, 0);
const cards = active.map((c) => { const cards = active.map((c) => {
const leg = legacy[c.provider] || {};
const legState = (leg.health && (leg.health.state || leg.health.status)) || null;
const rows = hidx[c.provider] || [];
const worst = rows.length ? worstBucket(rows.map((r) => bucketOf(r.state))) : (legState ? bucketOf(legState) : 'off');
const cred = c.credential || {}; const cred = c.credential || {};
const lastCall = (cidx[c.provider] || [])[0]; const observedRows = c.interfaces.filter((i) => i.observed);
const worst = observedRows.length ? worstBucket(observedRows.map((i) => bucketOf(i.observed_state))) : null;
const configured = cred.configured;
let bucket;
if (worst) bucket = worst;
else if (configured) bucket = 'plan';
else bucket = 'off';
const lastCall = cidx[c.provider];
if (!(c.provider in srcUiState)) srcUiState[c.provider] = { open: true, probing: false, probeResult: null }; if (!(c.provider in srcUiState)) srcUiState[c.provider] = { open: true, probing: false, probeResult: null };
const ifaceRows = c.interfaces.map((i) => { const groups = [];
const h = rows.find((r) => r.interface === i.interface); const groupMap = new Map();
return { c.interfaces.forEach((i) => {
api: i.interface, use: i.capability, const g = i.group || '其他';
ok: h ? (h.last_ok_at ? hm(h.last_ok_at) : '—') : '—', if (!groupMap.has(g)) groupMap.set(g, []);
lat: h ? h.last_latency_ms : null, groupMap.get(g).push(i);
st: h ? bucketOf(h.state) : 'off',
err: h && h.state !== 'ok' ? (h.last_error || h.last_fallback_reason || '') : '',
};
}); });
groupMap.forEach((items, name) => groups.push({ name, items }));
const ifaceRows = (group) => group.items.map((i) => ({
api: i.interface, use: i.capability,
ok: i.observed ? (i.observed_at ? hm(i.observed_at) : '—') : '—',
lat: i.observed ? i.observed_latency_ms : null,
st: i.observed ? bucketOf(i.observed_state) : (configured ? 'plan' : 'off'),
obs: i.observed, basis: i.observed_basis,
err: i.observed && i.observed_state !== 'ok' ? (i.observed_note || '') : '',
}));
return { return {
id: c.provider, name: c.label, role: c.role, tag: c.credential_key ? 'licensed' : 'free', id: c.provider, name: c.label, roleText: PROV_ROLE_TEXT[c.provider] || c.role,
bucket: worst, credential: cred, usedCount: ifaceRows.filter((r) => r.st !== 'off').length, totalCount: ifaceRows.length, tag: c.credential_key ? 'licensed' : 'free',
lastProbe: lastCall ? `${hm(lastCall.created_at)} · ${lastCall.ok ? `${lastCall.latency_ms}ms` : `× ${lastCall.error || 'ERROR'}`}` : '—', bucket, credential: cred,
latency: rows.length ? (rows.map((r) => r.last_latency_ms).filter((v) => v != null)[0] || null) : null, usedCount: c.interfaces.filter((i) => i.observed).length, totalCount: c.interfaces.length,
error: rows.find((r) => r.last_error)?.last_error || '', lastProbe: lastCall ? `${hm(lastCall.fetched_at)} ${lastCall.status === 'ok' ? `${lastCall.latency_ms ?? ''}ms` : `× ${lastCall.error || 'ERROR'}`}` : '—',
ifaceRows, latency: observedRows.length ? (observedRows.map((r) => r.observed_latency_ms).filter((v) => v != null)[0] || null) : null,
error: observedRows.find((r) => r.observed_state !== 'ok' && r.observed_note)?.observed_note || '',
groups: groups.map((g) => ({ name: g.name, rows: ifaceRows(g) })),
}; };
}); });
return { cards, reserved, totalIfaces, usedIfaces }; return { cards, reserved, totalIfaces, usedIfaces };
@@ -501,36 +566,48 @@ function sourceCardHtml(card) {
// 靠 note 字段区分,这里必须原样展示 note,不能一律折叠成"已配置"。 // 靠 note 字段区分,这里必须原样展示 note,不能一律折叠成"已配置"。
const credBadge = card.credential.note ? card.credential.note : const credBadge = card.credential.note ? card.credential.note :
card.credential.configured ? `已配置${card.credential.last4 ? ' · ' + card.credential.last4 : ''}` : '未配置'; card.credential.configured ? `已配置${card.credential.last4 ? ' · ' + card.credential.last4 : ''}` : '未配置';
const badges = [credBadge, card.role]; const badges = [credBadge];
if (card.credential.configured && card.credential.updated_at) badges.push(`${String(card.credential.updated_at).slice(5, 10).replace('-', '/')} 更新`);
badges.push(card.id === 'tushare' ? '主(无备)' : (card.id === 'ifind' ? '图表主源(未启用)' : (card.id === 'tencent' ? '实时层备' : '实时层主')));
let h = `<section class="panel" style="${ringStyle}">`; let h = `<section class="panel" style="${ringStyle}">`;
h += `<header class="src-head">` + h += `<header class="src-head">` +
ledHtml(ledSt, card.bucket === 'fail') + ledHtml(ledSt, card.bucket === 'fail') +
`<span class="src-name">${esc(card.name)}</span>` + `<span class="src-name">${esc(card.name)}</span>` +
`<span class="src-role">${esc(card.tag)}</span>` + `<span class="src-role">${esc(card.roleText)}</span>` +
pillHtml(card.bucket) + pillHtml(card.bucket, card.bucket === 'plan' ? '已配置 · 待观测' : '') +
`<span style="display:flex;gap:6px">${badges.map((b) => tagHtml(card.bucket === 'off' ? 'dashed' : 'dim', esc(b))).join('')}</span>` + `<span style="display:flex;gap:6px">${badges.map((b) => tagHtml(card.bucket === 'off' ? 'dashed' : 'dim', esc(b))).join('')}</span>` +
`<span class="flex1"></span>` + `<span class="flex1"></span>` +
`<button class="tbtn" data-probe="${card.id}"${ui.probing ? ' disabled' : ''}>${ui.probing ? '<span class="animate-spin" style="margin-right:4px">◌</span>探测中…' : '探测一次'}</button>` + `<button class="tbtn" data-probe="${card.id}"${ui.probing || card.bucket === 'off' ? ' disabled' : ''}>${ui.probing ? '<span class="animate-spin" style="margin-right:4px">◌</span>探测中…' : '探测一次'}</button>` +
`<button class="tbtn" data-toggle="${card.id}">${ui.open ? '收起接口清单' : '展开接口清单'}</button>` + (card.bucket === 'off'
? `<button class="tbtn" data-cred-hint="${card.id}">更新凭证</button>`
: `<button class="tbtn" data-toggle="${card.id}">${ui.open ? '收起接口清单' : '展开接口清单'}</button>`) +
`</header>`; `</header>`;
const probeCell = ui.probeResult
? `<span class="num" style="color:${ui.probeResult.ok ? '#34d399' : '#f87171'}">${ui.probeResult.t} ${ui.probeResult.ok ? `${ui.probeResult.ms}ms` : '× TIMEOUT'}</span>`
: `<span class="num"${card.lastProbe.includes('✓') ? ' style="color:#34d399"' : card.lastProbe.includes('×') ? ' style="color:#f87171"' : ''}>${esc(card.lastProbe)}</span>`;
h += `<div class="src-meta">` + h += `<div class="src-meta">` +
`<span>接口 <span class="num" style="color:#22d3ee">${card.usedCount}/${card.totalCount}</span> 在用</span>` + `<span>接口 <span class="num" style="color:${card.bucket === 'off' ? '#54637e' : '#22d3ee'}">${card.usedCount}/${card.totalCount}</span> 在用</span>` +
`<span class="msep">|</span><span>最后探测 <span class="num" style="color:${card.lastProbe.includes('×') ? '#f87171' : '#8b9bb4'}">${esc(card.lastProbe)}</span></span>` + `<span class="msep">|</span><span>最后探测 ${probeCell}</span>` +
`<span class="msep">|</span><span>延迟 ${latHtml(card.latency)}</span>` + `<span class="msep">|</span><span>延迟 ${latHtml(card.latency)}</span>` +
(card.error ? `<span class="msep">|</span><span style="color:#f87171">最近错误 ${esc(card.error)}</span>` : `<span class="msep">|</span><span style="color:#54637e">无最近错误</span>`) + (card.error ? `<span class="msep">|</span><span style="color:#f87171">最近错误 ${esc(card.error)}</span>` : (card.bucket !== 'off' ? `<span class="msep">|</span><span style="color:#54637e">无最近错误</span>` : '')) +
`</div>`; `</div>`;
if (ui.open && card.ifaceRows.length) { if (ui.open && card.bucket !== 'off') {
h += `<div class="mx-grid"><div style="min-width:0;grid-column:1/-1">` + h += `<div class="mx-grid">` + card.groups.map((g) => `<div style="min-width:0">` +
`<div class="lab" style="margin-bottom:6px">${esc(g.name)}</div>` +
`<table class="dtable"><thead><tr><th>接口</th><th>用途</th><th>最近成功</th><th>延迟</th><th>状态</th></tr></thead><tbody>` + `<table class="dtable"><thead><tr><th>接口</th><th>用途</th><th>最近成功</th><th>延迟</th><th>状态</th></tr></thead><tbody>` +
card.ifaceRows.map((it) => `<tr>` + g.rows.map((it) => `<tr>` +
`<td style="color:#22d3ee;font-size:11px;word-break:break-all;min-width:90px">${esc(it.api)}</td>` + `<td style="color:#22d3ee;font-size:11px;word-break:break-all;min-width:90px">${esc(it.api)}</td>` +
`<td style="color:#8b9bb4;font-size:11px;white-space:nowrap">${esc(it.use)}</td>` + `<td style="color:#8b9bb4;font-size:11px;white-space:nowrap">${esc(it.use)}</td>` +
`<td class="num" style="font-size:11px;color:${it.st === 'fail' ? '#f87171' : '#8b9bb4'}">${esc(it.ok)}</td>` + `<td class="num" style="font-size:11px;color:${it.ok.includes('×') || it.ok === '昨日' ? '#f87171' : '#8b9bb4'}">${esc(it.ok)}</td>` +
`<td style="font-size:11px">${latHtml(it.lat)}</td>` + `<td style="font-size:11px">${latHtml(it.lat)}</td>` +
`<td style="white-space:nowrap"><span style="display:flex;align-items:center;gap:6px">${pillHtml(it.st)}${it.err ? `<span style="font-size:10px;color:#f87171">${esc(it.err)}</span>` : ''}</span></td>` + `<td style="white-space:nowrap"><span style="display:flex;align-items:center;gap:6px">${pillHtml(it.st, it.st === 'plan' && it.obs === false ? '已配置 · 待观测' : '')}${it.err ? `<span style="font-size:10px;color:#f87171">${esc(it.err)}</span>` : ''}</span></td>` +
`</tr>`).join('') + `</tbody></table></div></div>`; `</tr>`).join('') +
} else if (!card.ifaceRows.length) { `</tbody></table></div>`).join('') + `</div>`;
h += `<div style="padding:16px"><div class="offbox">该来源尚无已登记接口</div></div>`; }
if (card.bucket === 'off') {
h += `<div style="padding:16px"><div class="offbox">凭证未配置 · ${card.totalCount} 个预留接口处于停用状态 · 配置后将成为图表 / 问财 / 实时快照主源` +
`<div style="margin-top:8px;display:flex;justify-content:center;gap:6px;flex-wrap:wrap">${card.groups[0].rows.map((i) => tagHtml('dashed', esc(i.api))).join('')}</div>` +
`</div></div>`;
} }
return h + `</section>`; return h + `</section>`;
} }
@@ -555,17 +632,57 @@ function sourcesHtml() {
function renderSourceCards() { function renderSourceCards() {
const el = $('srcCards'); const el = $('srcCards');
if (!el) return; if (!el) return;
const vm = buildSourcesVM(); el.innerHTML = buildSourcesVM().cards.map(sourceCardHtml).join('');
el.innerHTML = vm.cards.map(sourceCardHtml).join(''); paintSparks();
} }
/* ===================================================================== /* =====================================================================
数据血缘页 数据血缘页(接口级状态 + 已配置/待观测 + 真实频率 + 事件脉冲)
===================================================================== */ ===================================================================== */
function pageNameOf(consumer) { function pageNameOf(consumer) {
const idx = consumer.indexOf('('); const idx = consumer.indexOf('(');
return (idx > 0 ? consumer.slice(0, idx) : consumer).trim(); return (idx > 0 ? consumer.slice(0, idx) : consumer).trim();
} }
/* 一条血缘项的服务接口集合(primary/backup 冒号后的真实接口名) */
function servingInterfaces(item) {
const out = new Set([item.dataset]);
[item.primary_source, item.backup_source].forEach((src) => {
if (!src) return;
src.split(':').slice(1).join(':').split('+').forEach((api) => {
const name = api.trim();
if (name) out.add(name);
});
});
return out;
}
/* 接口级状态:只看本数据集真正走的那几条接口,不用提供商级汇总顶替。 */
function lineageItemRowsHealth(item) {
const apis = servingInterfaces(item);
return (item.live_provider_health || []).filter((r) => apis.has(r.interface));
}
function lineageItemBucket(item) {
const rows = lineageItemRowsHealth(item);
if (item.tier === 'official') {
if (item.publication) return item.publication.state === 'published' ? 'ok' : (item.publication.state === 'degraded' ? 'degrade' : 'fail');
return rows.length ? worstBucket(rows.map((r) => bucketOf(r.state))) : 'plan';
}
if (item.tier === 'licensed') {
const cred = ((RAW.catalog && RAW.catalog.items) || []).find((c) => c.provider === 'ifind');
if (cred && !(cred.credential && cred.credential.configured)) return 'off';
return rows.length ? worstBucket(rows.map((r) => bucketOf(r.state))) : 'off';
}
// provisional:已有数据来源;有观测行看接口级状态,暂无观测行=已配置/待观测
return rows.length ? worstBucket(rows.map((r) => bucketOf(r.state))) : 'plan';
}
function lineageItemNote(item) {
const rows = lineageItemRowsHealth(item);
const note = rows.find((r) => r.last_error)?.last_error || rows.find((r) => r.last_fallback_reason)?.last_fallback_reason || '';
if (note) return note;
if (item.tier === 'licensed') return 'iFinD 未配置 · 功能隐藏';
if (item.tier === 'provisional' && !rows.length) return '已配置 · 暂无观测记录';
return '';
}
const linFilter = { onlyBad: false, src: '全部', q: '' }; const linFilter = { onlyBad: false, src: '全部', q: '' };
let linHover = null; let linHover = null;
let GW = 1500, GH = 286; let GW = 1500, GH = 286;
@@ -579,18 +696,12 @@ function buildLineageVM() {
const srcState = {}; const srcState = {};
const dsState = {}; const dsState = {};
const pageState = {}; const pageState = {};
const pageDatasets = {};
const edgesSD = []; const edgesSD = [];
const edgesDP = []; const edgesDP = [];
const rows = []; const rows = [];
items.forEach((it) => { items.forEach((it) => {
const rowsHealth = it.live_provider_health || []; const rowsHealth = lineageItemRowsHealth(it);
let bucket; const bucket = lineageItemBucket(it);
if (it.tier === 'official') {
bucket = it.publication ? (it.publication.state === 'published' ? 'ok' : (it.publication.state === 'degraded' ? 'degrade' : 'fail')) : (rowsHealth.length ? worstBucket(rowsHealth.map((r) => bucketOf(r.state))) : 'plan');
} else {
bucket = rowsHealth.length ? worstBucket(rowsHealth.map((r) => bucketOf(r.state))) : 'off';
}
dsState[it.dataset] = bucket; dsState[it.dataset] = bucket;
(it.providers || []).forEach((p) => { (it.providers || []).forEach((p) => {
srcState[p] = worstBucket([srcState[p] || 'ok', ...rowsHealth.filter((r) => r.provider === p).map((r) => bucketOf(r.state))]); srcState[p] = worstBucket([srcState[p] || 'ok', ...rowsHealth.filter((r) => r.provider === p).map((r) => bucketOf(r.state))]);
@@ -601,16 +712,16 @@ function buildLineageVM() {
pages.forEach((pg) => { pages.forEach((pg) => {
pageState[pg] = worstBucket([pageState[pg] || 'ok', bucket]); pageState[pg] = worstBucket([pageState[pg] || 'ok', bucket]);
edgesDP.push([it.dataset, pg]); edgesDP.push([it.dataset, pg]);
(pageDatasets[pg] = pageDatasets[pg] || []).push(it);
}); });
const lat = rowsHealth.map((r) => r.last_latency_ms).find((v) => v != null);
const ok = it.tier === 'official' && it.publication ? it.publication.published_at : (rowsHealth.find((r) => r.last_ok_at) || {}).last_ok_at;
rows.push({ rows.push({
page: pages.join(' / ') || '—', item: consumers.join(' / '), dataset: it.dataset, page: pages.join(' / ') || '—', item: consumers.join(' / '), dataset: it.dataset,
via: `${it.primary_source}${it.backup_source ? ' → ' + it.backup_source : ''}`, via: `${it.primary_source}${it.backup_source ? ' → ' + it.backup_source : ''}`,
role: it.backup_source ? '主/备' : (it.tier === 'licensed' ? '主(授权)' : '主(无备)'), role: it.backup_source ? '主/备' : (it.tier === 'licensed' ? '主(授权)' : '主(无备)'),
ok: ok ? hm(ok) : '—', lat, st: bucket, freq: it.update_freq || '—',
note: rowsHealth.find((r) => r.last_error)?.last_error || rowsHealth.find((r) => r.last_fallback_reason)?.last_fallback_reason || (it.publication && it.publication.state !== 'published' ? '未正式发布' : ''), ok: it.tier === 'official' && it.publication ? it.publication.published_at : (rowsHealth.find((r) => r.last_ok_at) || {}).last_ok_at,
lat: rowsHealth.map((r) => r.last_latency_ms).find((v) => v != null),
st: bucket,
note: lineageItemNote(it),
}); });
}); });
const sources = Object.keys(srcState).map((id) => ({ id, name: PROV_LABEL[id] || id, st: srcState[id] })); const sources = Object.keys(srcState).map((id) => ({ id, name: PROV_LABEL[id] || id, st: srcState[id] }));
@@ -641,6 +752,7 @@ function lineageGraphHtml(vm) {
if (vm.datasets.length) h += `<text x="${gpos['d:' + vm.datasets[0].id].x}" y="12" text-anchor="middle" ${colStyle}>数据集 · ${vm.datasets.length}</text>`; if (vm.datasets.length) h += `<text x="${gpos['d:' + vm.datasets[0].id].x}" y="12" text-anchor="middle" ${colStyle}>数据集 · ${vm.datasets.length}</text>`;
if (vm.pages.length) h += `<text x="${gpos['p:' + vm.pages[0].id].x}" y="12" text-anchor="middle" ${colStyle}>主站页面 · ${vm.pages.length}</text>`; if (vm.pages.length) h += `<text x="${gpos['p:' + vm.pages[0].id].x}" y="12" text-anchor="middle" ${colStyle}>主站页面 · ${vm.pages.length}</text>`;
const dsBy = {}; vm.datasets.forEach((d) => { dsBy[d.id] = d.st; }); const dsBy = {}; vm.datasets.forEach((d) => { dsBy[d.id] = d.st; });
// 已连接线路全部有可辨识的定向流动(flowline),颜色按数据集真实状态
vm.edgesSD.forEach(([s, d]) => { vm.edgesSD.forEach(([s, d]) => {
if (!gpos['s:' + s] || !gpos['d:' + d]) return; if (!gpos['s:' + s] || !gpos['d:' + d]) return;
h += `<path data-link data-kind="sd" data-a="${esc(s)}" data-b="${esc(d)}" d="${linkPath(gpos['s:' + s], gpos['d:' + d])}" fill="none" stroke="${stColorOf(dsBy[d])}" stroke-width="1" stroke-opacity="0.28" class="${dsBy[d] === 'off' ? '' : 'flowline'}" style="transition:stroke-opacity .18s ease,stroke-width .18s ease"/>`; h += `<path data-link data-kind="sd" data-a="${esc(s)}" data-b="${esc(d)}" d="${linkPath(gpos['s:' + s], gpos['d:' + d])}" fill="none" stroke="${stColorOf(dsBy[d])}" stroke-width="1" stroke-opacity="0.28" class="${dsBy[d] === 'off' ? '' : 'flowline'}" style="transition:stroke-opacity .18s ease,stroke-width .18s ease"/>`;
@@ -716,35 +828,45 @@ function renderLinRows(vm) {
const rows = linRowsData(vm); const rows = linRowsData(vm);
if (cnt) cnt.textContent = `${rows.length}/${vm.rows.length}`; if (cnt) cnt.textContent = `${rows.length}/${vm.rows.length}`;
const cls = (st) => st === 'fail' ? 'row-fail' : st === 'degrade' || st === 'slow' ? 'row-slow' : st === 'off' ? 'row-off' : ''; const cls = (st) => st === 'fail' ? 'row-fail' : st === 'degrade' || st === 'slow' ? 'row-slow' : st === 'off' ? 'row-off' : '';
const stLabel = (r) => (r.st === 'plan' && r.note.includes('待观测')) ? '已配置 · 待观测' : '';
body.innerHTML = rows.map((r) => `<tr class="${cls(r.st)}">` + body.innerHTML = rows.map((r) => `<tr class="${cls(r.st)}">` +
`<td style="color:#e8f1ff;white-space:nowrap">${esc(r.page)}</td>` + `<td style="color:#e8f1ff;white-space:nowrap">${esc(r.page)}</td>` +
`<td style="color:#8b9bb4;font-size:11px;white-space:nowrap">${esc(r.item)}</td>` + `<td style="color:#8b9bb4;font-size:11px;white-space:nowrap">${esc(r.item)}</td>` +
`<td style="color:#22d3ee;font-size:11px">${esc(r.dataset)}</td>` + `<td style="color:#22d3ee;font-size:11px">${esc(r.dataset)}</td>` +
`<td style="color:#8b9bb4;font-size:11px;max-width:260px;overflow:hidden;text-overflow:ellipsis;white-space:nowrap" title="${esc(r.via)}">${esc(r.via)}</td>` + `<td style="color:#8b9bb4;font-size:11px;max-width:260px;overflow:hidden;text-overflow:ellipsis;white-space:nowrap" title="${esc(r.via)}">${esc(r.via)}</td>` +
`<td style="color:#54637e;font-size:11px;white-space:nowrap">${esc(r.role)}</td>` + `<td style="color:#54637e;font-size:11px;white-space:nowrap">${esc(r.role)}</td>` +
`<td class="num" style="font-size:11px;white-space:nowrap;color:#8b9bb4">${esc(r.ok)}</td>` + `<td class="num" style="color:#8b9bb4;font-size:11px;white-space:nowrap">${esc(r.freq)}</td>` +
`<td>${latHtml(r.lat)}</td>` + `<td class="num" style="font-size:11px;white-space:nowrap;color:#8b9bb4">${r.ok ? esc(hm(r.ok)) : '—'}</td>` +
`<td><span style="display:flex;align-items:center;gap:8px;white-space:nowrap">${pillHtml(r.st)}${r.note ? `<span style="font-size:10px;color:#54637e">${esc(r.note)}</span>` : ''}</span></td>` + `<td class="num" style="font-size:11px;color:${r.lat != null && r.lat > 1000 ? '#fbbf24' : '#8b9bb4'}">${r.lat != null ? Number(r.lat).toLocaleString('en-US') + 'ms' : ''}</td>` +
`</tr>`).join('') || `<tr><td colspan="8" class="empty-hint">无匹配记录 · 调整筛选条件</td></tr>`; `<td><span style="display:flex;align-items:center;gap:8px;white-space:nowrap">${pillHtml(r.st, stLabel(r))}${r.note && !stLabel(r) ? `<span style="font-size:10px;color:#54637e">${esc(r.note)}</span>` : ''}</span></td>` +
`</tr>`).join('') || `<tr><td colspan="9" class="empty-hint">无匹配记录 · 调整筛选条件</td></tr>`;
if (bodyN) bodyN.innerHTML = rows.map((r) => `<tr class="${cls(r.st)}">` + if (bodyN) bodyN.innerHTML = rows.map((r) => `<tr class="${cls(r.st)}">` +
`<td style="color:#e8f1ff;white-space:nowrap">${esc(r.page)}<span class="sub" style="color:#8b9bb4">${esc(r.item)}</span></td>` + `<td style="color:#e8f1ff;white-space:nowrap">${esc(r.page)}<span class="sub" style="color:#8b9bb4">${esc(r.item)}</span></td>` +
`<td style="color:#22d3ee;font-size:11px">${esc(r.dataset)}</td>` + `<td style="color:#22d3ee;font-size:11px">${esc(r.dataset)}</td>` +
`<td style="color:#8b9bb4;font-size:11px;word-break:break-all">${esc(r.via)}</td>` + `<td style="color:#8b9bb4;font-size:11px;word-break:break-all">${esc(r.via)}</td>` +
`<td class="num" style="font-size:11px;white-space:nowrap;color:#8b9bb4">${esc(r.ok)}<span class="sub">${r.lat != null ? r.lat + 'ms' : '—'}</span></td>` + `<td class="num" style="color:#8b9bb4;font-size:11px;white-space:nowrap">${esc(r.freq)}<span class="sub">${esc(r.role)}</span></td>` +
`<td>${pillHtml(r.st)}${r.note ? `<span class="sub">${esc(r.note)}</span>` : ''}</td>` + `<td class="num" style="font-size:11px;white-space:nowrap;color:#8b9bb4">${r.ok ? esc(hm(r.ok)) : '—'}<span class="sub">${r.lat != null ? r.lat + 'ms' : ''}</span></td>` +
`</tr>`).join('') || `<tr><td colspan="5" class="empty-hint">无匹配记录 · 调整筛选条件</td></tr>`; `<td>${pillHtml(r.st, stLabel(r))}${r.note && !stLabel(r) ? `<span class="sub">${esc(r.note)}</span>` : ''}</td>` +
`</tr>`).join('') || `<tr><td colspan="6" class="empty-hint">无匹配记录 · 调整筛选条件</td></tr>`;
} }
function lineageHtml() { function lineageHtml() {
const vm = buildLineageVM(); const vm = buildLineageVM();
const anomalyPages = new Set(); const failDs = vm.datasets.filter((d) => d.st === 'fail');
vm.rows.forEach((r) => { if (r.st === 'fail') (r.page || '').split(' / ').forEach((p) => anomalyPages.add(p)); }); let impactTag;
if (failDs.length) {
const pages = new Set();
failDs.forEach((d) => { (adj.dp.get(d.id) || new Set()).forEach((p) => pages.add(p)); });
impactTag = tagHtml('rd', `${failDs.map((d) => d.id).join('/')} 断链影响 ${pages.size}`);
} else {
impactTag = tagHtml('mint', '暂无断链影响');
}
const html = `<div class="col">` + const html = `<div class="col">` +
`<div class="panel strip" style="gap:8px 20px">` + `<div class="panel strip" style="gap:8px 20px">` +
`<span style="font-size:13px;color:#e8f1ff">血缘地图</span>` + `<span style="font-size:13px;color:#e8f1ff">血缘地图</span>` +
`<span class="num" style="font-size:11px;color:#8b9bb4"><span style="color:#22d3ee">${vm.sources.length}</span> 源 → <span style="color:#22d3ee">${vm.datasets.length}</span> 数据集 → <span style="color:#22d3ee">${vm.pages.length}</span> 页面</span>` + `<span class="num" style="font-size:11px;color:#8b9bb4"><span style="color:#22d3ee">${vm.sources.length}</span> 源 → <span style="color:#22d3ee">${vm.datasets.length}</span> 数据集 → <span style="color:#22d3ee">${vm.pages.length}</span> 页面</span>` +
`<span class="legend" style="color:#54637e;letter-spacing:0"><span>${ledHtml('ok')}健康</span><span>${ledHtml('slow')}慢/降级</span><span>${ledHtml('fail')}失败</span><span>${ledHtml('plan')}计划</span><span>${ledHtml('off')}未配置</span></span>` + `<span class="legend" style="color:#54637e;letter-spacing:0"><span>${ledHtml('ok')}健康</span><span>${ledHtml('slow')}慢/降级</span><span>${ledHtml('fail')}失败</span><span>${ledHtml('plan')}计划</span><span>${ledHtml('off')}未配置</span></span>` +
`<span class="flex1"></span>` + `<span class="flex1"></span>` +
(anomalyPages.size ? tagHtml('rd', `${[...anomalyPages][0]} 断链影响 ${anomalyPages.size}`) : tagHtml('mint', '暂无断链影响')) + impactTag +
`</div>` + `</div>` +
`<section class="panel"><div style="position:relative">${lineageGraphHtml(vm)}</div></section>` + `<section class="panel"><div style="position:relative">${lineageGraphHtml(vm)}</div></section>` +
`<section class="panel">` + `<section class="panel">` +
@@ -755,11 +877,10 @@ function lineageHtml() {
`<span class="num" style="font-size:10px;color:#54637e" id="linCount"></span>` + `<span class="num" style="font-size:10px;color:#54637e" id="linCount"></span>` +
`<input class="lin-q" id="linQ" placeholder="搜索页面 / 接口 / 数据集" value="${esc(linFilter.q)}">` + `<input class="lin-q" id="linQ" placeholder="搜索页面 / 接口 / 数据集" value="${esc(linFilter.q)}">` +
`</div>` + `</div>` +
`<div class="lin-scroll"><table class="dtable lin-table"><thead><tr><th>主站页面</th><th>数据项</th><th>数据集</th><th>源</th><th>主/备</th><th>最近成功</th><th>延迟</th><th>状态</th></tr></thead><tbody id="linBody"></tbody></table>` + `<div class="lin-scroll"><table class="dtable lin-table"><thead><tr><th>主站页面</th><th>数据项</th><th>数据集</th><th>源 : 接口</th><th>主/备</th><th>更新频率</th><th>最近成功</th><th>延迟</th><th>状态</th></tr></thead><tbody id="linBody"></tbody></table>` +
`<table class="dtable lin-table-n"><thead><tr><th>主站页面</th><th>数据集</th><th>源</th><th>最近成功</th><th>状态</th></tr></thead><tbody id="linBodyN"></tbody></table></div>` + `<table class="dtable lin-table-n"><thead><tr><th>主站页面</th><th>数据集</th><th>源 : 接口</th><th>更新频率</th><th>最近成功</th><th>状态</th></tr></thead><tbody id="linBodyN"></tbody></table></div>` +
`</section>` + `</section>` +
`</div>`; `</div>`;
window.__linVM = vm;
return html; return html;
} }
@@ -775,16 +896,19 @@ function renderPage(p) {
state.page = p; state.page = p;
mainEl.dataset.page = p; mainEl.dataset.page = p;
document.querySelectorAll('[data-nav]').forEach((b) => b.classList.toggle('act', b.dataset.nav === p)); document.querySelectorAll('[data-nav]').forEach((b) => b.classList.toggle('act', b.dataset.nav === p));
if (p === 'overview') mainEl.innerHTML = overviewHtml(); if (p === 'overview') { mainEl.innerHTML = overviewHtml(); paintSparks(); flashChangedLatencies(); }
else if (p === 'sources') mainEl.innerHTML = sourcesHtml(); else if (p === 'sources') { mainEl.innerHTML = sourcesHtml(); paintSparks(); }
else { else {
mainEl.innerHTML = lineageHtml(); mainEl.innerHTML = lineageHtml();
renderLinRows(window.__linVM); const vm = buildLineageVM();
applyLinHover(window.__linVM); window.__linVM = vm;
renderLinRows(vm);
applyLinHover(vm);
const svg = $('linSvg'); const svg = $('linSvg');
if (svg) svg.addEventListener('mouseleave', () => setLinHover(null, window.__linVM)); if (svg) svg.addEventListener('mouseleave', () => setLinHover(null, window.__linVM));
fitLineageTable(); fitLineageTable();
} }
window.scrollTo(0, 0);
updatePhaseTag(); updatePhaseTag();
} }
window.addEventListener('hashchange', () => renderPage(pageFromHash())); window.addEventListener('hashchange', () => renderPage(pageFromHash()));
@@ -795,11 +919,10 @@ window.addEventListener('resize', () => {
resizeT = setTimeout(() => { if (state.page === 'lineage') renderPage('lineage'); }, 180); resizeT = setTimeout(() => { if (state.page === 'lineage') renderPage('lineage'); }, 180);
}); });
/* HEL-529: 血缘详情表随真实数据集数量(当前 17 条,非打样假设的 13 条) /* HEL-529: 血缘详情表随真实数据集数量动态撑高;宽档(>1100px,覆盖
动态撑高;宽档(>1100px,覆盖 1920×920/1440×900 验收位)给表格容器一个 1920×920/1440×900 验收位)给表格容器一个基于剩余可视高度计算的
基于剩余可视高度计算的 max-height + 内部滚动,保证整页 scrollHeight max-height + 内部滚动,保证整页 scrollHeight 始终 <= clientHeight
始终 <= clientHeight(硬性验收项),同时全部行仍可在容器内滚动查看。 同时全部行仍可在容器内滚动查看。窄档(<=1100px)整页纵滚。 */
窄档(<=1100px)保持原有整页纵向滚动,不设上限。 */
function fitLineageTable() { function fitLineageTable() {
const wrap = document.querySelector('.lin-scroll'); const wrap = document.querySelector('.lin-scroll');
const tape = document.querySelector('.tape'); const tape = document.querySelector('.tape');
@@ -816,6 +939,8 @@ function fitLineageTable() {
function updatePhaseTag() { function updatePhaseTag() {
const ov = RAW.overview; const ov = RAW.overview;
if (!ov) return; if (!ov) return;
const bad = (ov.anomalies || []).length > 0;
$('radarBlipBad').style.display = bad ? '' : 'none';
$('phaseTag').innerHTML = `${ledHtml('rev', true)} ${esc(ov.session_phase || '')}`; $('phaseTag').innerHTML = `${ledHtml('rev', true)} ${esc(ov.session_phase || '')}`;
} }
@@ -824,6 +949,11 @@ mainEl.addEventListener('click', (e) => {
if (pb) { doProbe(pb.dataset.probe); return; } if (pb) { doProbe(pb.dataset.probe); return; }
const tg = e.target.closest('[data-toggle]'); const tg = e.target.closest('[data-toggle]');
if (tg) { const ui = srcUiState[tg.dataset.toggle]; if (ui) { ui.open = !ui.open; renderSourceCards(); } return; } if (tg) { const ui = srcUiState[tg.dataset.toggle]; if (ui) { ui.open = !ui.open; renderSourceCards(); } return; }
const ch = e.target.closest('[data-cred-hint]');
if (ch) {
toast('iFinD 凭证需在部署侧配置(DATAHUB_IFIND_* 环境变量或 /v1/credentials/ifind),后台暂未开放在线写入', 'err');
return;
}
const goto = e.target.closest('[data-goto-lineage]'); const goto = e.target.closest('[data-goto-lineage]');
if (goto) { linFilter.q = ''; location.hash = 'lineage'; return; } if (goto) { linFilter.q = ''; location.hash = 'lineage'; return; }
const inc = e.target.closest('[data-incident-action]'); const inc = e.target.closest('[data-incident-action]');
@@ -846,14 +976,35 @@ async function doProbe(id) {
if (!ui || ui.probing) return; if (!ui || ui.probing) return;
ui.probing = true; ui.probing = true;
renderSourceCards(); renderSourceCards();
try { await api(`/admin/api/sources/${id}/probe`, { method: 'POST', body: '{}' }); } catch (err) { toast(err.message, 'err'); } try {
await api(`/admin/api/sources/${id}/probe`, { method: 'POST', body: '{}' });
ui.probeResult = { t: `${p2(new Date().getHours())}:${p2(new Date().getMinutes())}`, ok: true, ms: null };
} catch (err) {
ui.probeResult = { t: `${p2(new Date().getHours())}:${p2(new Date().getMinutes())}`, ok: false, ms: null };
toast(err.message, 'err');
}
ui.probing = false; ui.probing = false;
await refreshQuiet(); await refreshQuiet();
renderSourceCards(); renderSourceCards();
} }
/* ---------- EVENT TAPE真实事件差量渲染 ---------- */ /* ---------- EVENT TAPE预装最近真实事件 + 持续滚动 + 新事件闪动 ---------- */
const tapeSeen = { call: 0, provCall: 0, audit: 0, job: 0, first: true }; const tapeSeen = { call: 0, provCall: 0, audit: 0, job: 0 };
let tapeBuf = [];
let tapeSeeded = false;
function tapeEventHtml(e, fresh) {
return `<span class="tape-item${fresh ? ' flash' : ''}"><span class="tt num">${esc(hm(e.ts))}</span><span class="ts">${esc(e.src)}</span><span class="tm">${esc(e.msg)}</span>` +
(e.ok ? `<span class="tok num">✓${e.ms ? ` ${e.ms}ms` : ''}</span>` : '<span class="tbad">× 失败</span>') +
`<span class="tsep">///</span></span>`;
}
function renderTape(freshIds) {
const items = tapeBuf.slice(-14);
const html = items.length
? items.map((e) => tapeEventHtml(e, freshIds && freshIds.has(e.sortId))).join('')
: `<span class="tape-item"><span class="tm" style="color:#54637e">暂无真实事件 · 等待首次调用/任务</span><span class="tsep">///</span></span>`;
$('tapeA').innerHTML = html;
$('tapeB').innerHTML = html;
}
function tapeEventsFromRaw() { function tapeEventsFromRaw() {
const evs = []; const evs = [];
((RAW.overview && RAW.overview.recent_calls) || []).forEach((c) => { if (c.id > tapeSeen.call) evs.push({ ts: c.created_at, src: 'tushare', msg: c.endpoint, ok: c.ok, ms: c.latency_ms, sortId: 'c' + c.id }); }); ((RAW.overview && RAW.overview.recent_calls) || []).forEach((c) => { if (c.id > tapeSeen.call) evs.push({ ts: c.created_at, src: 'tushare', msg: c.endpoint, ok: c.ok, ms: c.latency_ms, sortId: 'c' + c.id }); });
@@ -867,16 +1018,25 @@ function tapeEventsFromRaw() {
tapeSeen.job = Math.max(tapeSeen.job, maxOf((RAW.jobs && RAW.jobs.runs) || [], 'id')); tapeSeen.job = Math.max(tapeSeen.job, maxOf((RAW.jobs && RAW.jobs.runs) || [], 'id'));
return evs; return evs;
} }
let tapeBuf = []; function seedTape() {
// 首次加载:预装最近的真实事件(打样 EVENT TAPE 开场即有内容滚动)
const evs = [];
((RAW.providers && RAW.providers.recent_calls) || []).forEach((c) => evs.push({ ts: c.fetched_at, src: c.provider, msg: c.interface, ok: c.status === 'ok', ms: c.latency_ms, sortId: 'p' + c.id }));
((RAW.overview && RAW.overview.recent_calls) || []).forEach((c) => evs.push({ ts: c.created_at, src: 'tushare', msg: c.endpoint, ok: c.ok, ms: c.latency_ms, sortId: 'c' + c.id }));
((RAW.jobs && RAW.jobs.runs) || []).forEach((r) => evs.push({ ts: r.finished_at || r.started_at, src: 'scheduler', msg: `${r.job_id} ${r.state}`, ok: r.state !== 'failed', sortId: 'j' + r.id }));
((RAW.audit && RAW.audit.items) || []).forEach((a) => evs.push({ ts: a.created_at, src: 'hub_admin', msg: `${a.action} ${a.target || ''}`.trim(), ok: true, sortId: 'a' + a.id }));
evs.sort((x, y) => String(x.ts).localeCompare(String(y.ts)));
tapeBuf = evs.slice(-14);
tapeSeeded = true;
renderTape(null);
}
function updateTape() { function updateTape() {
if (!tapeSeeded) { seedTape(); return; }
const evs = tapeEventsFromRaw(); const evs = tapeEventsFromRaw();
if (tapeSeen.first) { tapeSeen.first = false; return; } // 首次加载只建基线,不补播历史事件
if (!evs.length) return; if (!evs.length) return;
const fresh = new Set(evs.map((e) => e.sortId));
tapeBuf = tapeBuf.concat(evs).slice(-16); tapeBuf = tapeBuf.concat(evs).slice(-16);
const html = tapeBuf.map((e) => `<span class="tape-item"><span class="tt num">${esc(hm(e.ts))}</span><span class="ts">${esc(e.src)}</span><span class="tm">${esc(e.msg)}</span>` + renderTape(fresh);
(e.ok ? `<span class="tok num">✓${e.ms ? ` ${e.ms}ms` : ''}</span>` : '<span class="tbad">× 失败</span>') +
`<span class="tsep">///</span></span>`).join('');
$('tapeA').innerHTML = html; $('tapeB').innerHTML = html;
} }
/* ---------- 抽屉:调度任务 / 盘后发布 / 审计 ---------- */ /* ---------- 抽屉:调度任务 / 盘后发布 / 审计 ---------- */
@@ -1011,18 +1171,22 @@ function toast(msg, kind) {
toastT = setTimeout(() => { el.classList.remove('show'); }, 3200); toastT = setTimeout(() => { el.classList.remove('show'); }, 3200);
} }
/* ================= 顶栏:时钟 / 减少动态 / 主题 ================= */ /* ================= 顶栏:时钟 / 减少动态 ================= */
let clockEls = null;
function updateClock() { function updateClock() {
if (!clockEls) clockEls = { d: $('ckD'), h: $('ckH'), m: $('ckM'), s: $('ckS') };
const now = new Date(); const now = new Date();
$('ckD').textContent = `${now.getFullYear()}-${p2(now.getMonth() + 1)}-${p2(now.getDate())}`; const ds = `${now.getFullYear()}-${p2(now.getMonth() + 1)}-${p2(now.getDate())}`;
$('ckH').textContent = p2(now.getHours()); if (clockEls.d.textContent !== ds) clockEls.d.textContent = ds;
$('ckM').textContent = p2(now.getMinutes()); const hs = p2(now.getHours()), ms = p2(now.getMinutes()), ss = p2(now.getSeconds());
$('ckS').textContent = p2(now.getSeconds()); if (clockEls.h.textContent !== hs) clockEls.h.textContent = hs;
if (clockEls.m.textContent !== ms) clockEls.m.textContent = ms;
if (clockEls.s.textContent !== ss) clockEls.s.textContent = ss;
} }
function applyCalm(on) { function applyCalm(on) {
state.calm = on; state.calm = on;
// 挂在 #app(覆盖 appRoot 及其兄弟 drawer/modal/toast),不能只挂 #appRoot // 挂在 body + #app(覆盖 appRoot 及其兄弟 drawer/modal/toast),否则
// 否则 .reduce-motion * 与 .reduce-motion .drawer/.toast 选不中旁路层过渡。 // .reduce-motion * 选不中旁路层过渡。
document.body.classList.toggle('reduce-motion', on); document.body.classList.toggle('reduce-motion', on);
const app = $('app'); const app = $('app');
if (app) app.classList.toggle('reduce-motion', on); if (app) app.classList.toggle('reduce-motion', on);
@@ -1037,7 +1201,62 @@ if (window.matchMedia) {
mq.addEventListener ? mq.addEventListener('change', syncMq) : mq.addListener(syncMq); mq.addEventListener ? mq.addEventListener('change', syncMq) : mq.addListener(syncMq);
} }
/* ================= 动态更新:200ms tick(打样节奏) =================
- 时钟每秒翻字(冒号 blink 由 CSS 承担)
- 缓存年龄依据真实时间戳每秒增长(>30s 变琥珀)
- 倒计时按真实 next_retry_at 每秒递减
- sparkline 端点闪烁由 CSS 承担,折线只在真实历史变化时重画
- 延迟数字只在真实值变化时闪一下(flash),绝不随机伪造数值 */
function paintSparks() {
document.querySelectorAll('svg.spark').forEach((svg) => sparkPaint(svg));
}
const prevLat = new Map();
function flashChangedLatencies() {
document.querySelectorAll('[data-lat]').forEach((el) => {
const jid = el.closest('[data-jid]') ? el.closest('[data-jid]').dataset.jid : (el.closest('tr') ? 'row:' + [...el.closest('tr').children].indexOf(el) : '');
const v = el.dataset.lat;
const key = state.page + ':' + (el.closest('[data-dataset]')?.dataset.dataset || jid) + ':' + v;
const prev = prevLat.get(el);
if (prev !== undefined && prev !== v) {
el.classList.remove('flash');
void el.offsetWidth;
el.classList.add('flash');
}
prevLat.set(el, v);
});
}
function updateDynamics() {
const now = Date.now();
document.querySelectorAll('[data-age-epoch]').forEach((el) => {
const ms = +el.dataset.ageEpoch;
if (!ms || Number.isNaN(ms)) return;
const s = Math.max(0, Math.floor((now - ms) / 1000));
const txt = fmtAgeSec(s);
if (el.textContent !== txt) el.textContent = txt;
const color = s > 30 ? '#fbbf24' : '#8b9bb4';
if (el.style.color !== color) el.style.color = color;
});
document.querySelectorAll('[data-cd-epoch]').forEach((el) => {
const target = +el.dataset.cdEpoch;
if (!target || Number.isNaN(target)) return;
const left = Math.floor((target - now) / 1000);
if (left > 0) {
const m = Math.floor(left / 60), s = left % 60;
const txt = el.dataset.cdfmt === 'ms' ? `${m}m${p2(s)}s` : `${p2(m)}:${p2(s)}`;
if (el.textContent !== txt) el.textContent = txt;
} else if (el.textContent !== '已到点') {
el.textContent = '已到点';
}
});
}
/* ================= 轮询:拉取真实数据、驱动渲染、可见性/网络门禁 ================= */ /* ================= 轮询:拉取真实数据、驱动渲染、可见性/网络门禁 ================= */
const dataVersion = { overview: '', sources: '', lineage: '' };
function snapshotVersion(part) {
if (part === 'overview') return JSON.stringify([RAW.overview, RAW.batches && RAW.batches.batches, RAW.providers && RAW.providers.health, RAW.catalog]);
if (part === 'sources') return JSON.stringify([RAW.catalog, RAW.providers && RAW.providers.recent_calls, RAW.sources]);
return JSON.stringify(RAW.lineage);
}
const Poller = (() => { const Poller = (() => {
let timer = null; let timer = null;
let running = false; let running = false;
@@ -1047,9 +1266,32 @@ const Poller = (() => {
try { try {
await fetchAll(); await fetchAll();
updateTape(); updateTape();
if (state.page === 'overview') mainEl.innerHTML = overviewHtml(); const vOv = snapshotVersion('overview'), vSrc = snapshotVersion('sources'), vLin = snapshotVersion('lineage');
else if (state.page === 'sources') renderSourceCards(); if (state.page === 'overview' && vOv !== dataVersion.overview) {
else if (state.page === 'lineage') { /* 血缘图交互态复杂,仅在用户操作或下次轮询后整体重绘一次 */ mainEl.innerHTML = lineageHtml(); renderLinRows(window.__linVM); applyLinHover(window.__linVM); fitLineageTable(); } dataVersion.overview = vOv;
mainEl.innerHTML = overviewHtml();
paintSparks();
flashChangedLatencies();
} else if (state.page === 'overview') {
flashChangedLatencies();
}
if (state.page === 'sources' && vSrc !== dataVersion.sources) {
dataVersion.sources = vSrc;
renderSourceCards();
}
if (state.page === 'lineage' && vLin !== dataVersion.lineage) {
dataVersion.lineage = vLin;
mainEl.innerHTML = lineageHtml();
const vm = buildLineageVM();
window.__linVM = vm;
renderLinRows(vm);
applyLinHover(vm);
const svg = $('linSvg');
if (svg) svg.addEventListener('mouseleave', () => setLinHover(null, window.__linVM));
fitLineageTable();
}
// 真实调用事件脉冲:血缘页上对应来源节点亮一圈 ledPing 节奏的光环
pulseSourcesOnNewCalls();
updatePhaseTag(); updatePhaseTag();
setNetOk(true); setNetOk(true);
} catch (err) { } catch (err) {
@@ -1072,10 +1314,39 @@ const Poller = (() => {
}; };
})(); })();
/* 真实调用脉冲:provider_call_log 出现新调用时,给血缘图对应来源节点
叠一圈与 ledPing 同节奏(1.6s ease-out)的扩散光环;这是事件动效,
与 gnode-pulse 的静态呼吸不同,一次即消散。 */
let pulseSeenCall = 0;
function pulseSourcesOnNewCalls() {
if (state.page !== 'lineage') {
pulseSeenCall = ((RAW.providers && RAW.providers.recent_calls) || []).reduce((m, c) => Math.max(m, c.id || 0), pulseSeenCall);
return;
}
const calls = (RAW.providers && RAW.providers.recent_calls) || [];
const fresh = calls.filter((c) => (c.id || 0) > pulseSeenCall);
pulseSeenCall = calls.reduce((m, c) => Math.max(m, c.id || 0), pulseSeenCall);
fresh.forEach((c) => {
const g = document.querySelector(`#linSvg g[data-nid="${CSS.escape(c.provider)}"]`);
if (!g) return;
const circle = g.querySelector('circle');
if (!circle) return;
const ring = document.createElementNS('http://www.w3.org/2000/svg', 'circle');
ring.setAttribute('cx', circle.getAttribute('cx'));
ring.setAttribute('cy', circle.getAttribute('cy'));
ring.setAttribute('r', '5');
ring.setAttribute('fill', 'none');
ring.setAttribute('stroke', '#22d3ee');
ring.setAttribute('stroke-width', '1.2');
ring.setAttribute('class', 'node-ping');
g.appendChild(ring);
setTimeout(() => ring.remove(), 1700);
});
}
function resetTapeBaselines() { function resetTapeBaselines() {
// 恢复轮询时静默重建基线,不补播暂停期间错过的事件。 // 恢复轮询时静默重建基线,不补播暂停期间错过的事件(但保留已上带内容)
tapeSeen.call = tapeSeen.provCall = tapeSeen.audit = tapeSeen.job = Number.MAX_SAFE_INTEGER; tapeSeen.call = tapeSeen.provCall = tapeSeen.audit = tapeSeen.job = Number.MAX_SAFE_INTEGER;
tapeSeen.first = true;
} }
document.addEventListener('visibilitychange', () => { document.addEventListener('visibilitychange', () => {
if (document.hidden) { Poller.stop(); } if (document.hidden) { Poller.stop(); }
@@ -1090,11 +1361,15 @@ async function refreshQuiet() {
/* ================= 启动 ================= */ /* ================= 启动 ================= */
async function startApp() { async function startApp() {
setInterval(updateClock, 1000);
updateClock();
await fetchAll(); await fetchAll();
tapeSeen.first = true; dataVersion.overview = snapshotVersion('overview');
dataVersion.sources = snapshotVersion('sources');
dataVersion.lineage = snapshotVersion('lineage');
pulseSeenCall = ((RAW.providers && RAW.providers.recent_calls) || []).reduce((m, c) => Math.max(m, c.id || 0), 0);
renderPage(pageFromHash()); renderPage(pageFromHash());
updateTape(); // 预装最近真实事件并开始滚动
updateClock();
setInterval(() => { updateClock(); updateDynamics(); }, 200); // 打样 tick 节奏
await Poller.start(); await Poller.start();
} }
+5 -3
View File
@@ -46,8 +46,10 @@
<div class="ring1" style="inset:5.72px"></div> <div class="ring1" style="inset:5.72px"></div>
<div class="ring2" style="inset:9.88px"></div> <div class="ring2" style="inset:9.88px"></div>
<div class="center"></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>
<div class="brand"><span class="b1">小白复盘 <span style="color:#22d3ee">·</span> 数据中枢</span><span class="b2">DATA-HUB · HEL-529</span></div> <div class="brand"><span class="b1">小白复盘 <span style="color:#22d3ee">·</span> 数据中枢</span><span class="b2">DATA-HUB</span></div>
<span class="vsep"></span> <span class="vsep"></span>
<nav class="nav" id="navEl"> <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="overview"><span class="tri"></span>运行总览<span class="en">OVERVIEW</span></button>
@@ -57,7 +59,7 @@
<span class="flex1"></span> <span class="flex1"></span>
<span class="mdtag" id="phaseTag"></span> <span class="mdtag" id="phaseTag"></span>
<span class="livespan" id="liveSpan"><span class="livedot pulse-dot"></span>LIVE</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="csep">:</span><span id="ckM"></span><span class="csep">:</span><span id="ckS"></span></span></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> <span class="vsep"></span>
<button class="tbtn" id="opsBtn" style="font-size:10px">调度 / 发布 / 审计</button> <button class="tbtn" id="opsBtn" style="font-size:10px">调度 / 发布 / 审计</button>
<button class="tbtn" id="calmBtn" style="font-size:10px">减少动态</button> <button class="tbtn" id="calmBtn" style="font-size:10px">减少动态</button>
@@ -69,7 +71,7 @@
<main class="main" id="mainEl"></main> <main class="main" id="mainEl"></main>
<footer class="tape"> <footer class="tape">
<span class="tape-label" id="tapeLed">EVENT TAPE</span> <span class="tape-label" id="tapeLed"><span class="led rev pulse"></span>EVENT TAPE</span>
<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> <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>
</footer> </footer>
</div> </div>
+8
View File
@@ -367,6 +367,14 @@ body {
.lin-hint { position: absolute; left: 24px; bottom: 4px; font-size: 10px; color: #54637e; letter-spacing: .1em; } .lin-hint { position: absolute; left: 24px; bottom: 4px; font-size: 10px; color: #54637e; letter-spacing: .1em; }
@keyframes gnodePulse { 0%,100% { opacity: 1; } 50% { opacity: .55; } } @keyframes gnodePulse { 0%,100% { opacity: 1; } 50% { opacity: .55; } }
.gnode-pulse { animation: gnodePulse 1.8s ease-in-out infinite; } .gnode-pulse { animation: gnodePulse 1.8s ease-in-out infinite; }
/* 真实调用事件脉冲(HEL-529 数据修正第 6 项):血缘图来源节点在真实调用
到达时叠一圈与 ledPing 同节奏(1.6s ease-out 放大淡出)的光环;
事件动效,一次即消散,与 gnode-pulse 的静态呼吸区分。 */
.node-ping { transform-box: fill-box; transform-origin: center; animation: nodePing 1.6s ease-out forwards; }
@keyframes nodePing {
0% { transform: scale(.4); opacity: .8; }
80%, 100% { transform: scale(3.2); opacity: 0; }
}
.lin-filter { display: flex; flex-wrap: wrap; align-items: center; gap: 8px; padding: 10px 12px; border-bottom: 1px solid #1a2540; } .lin-filter { display: flex; flex-wrap: wrap; align-items: center; gap: 8px; padding: 10px 12px; border-bottom: 1px solid #1a2540; }
.lin-sel, .lin-q { background: #0c1220; border: 1px solid #243152; border-radius: 4px; font-size: 11px; color: #8b9bb4; padding: 6px 8px; outline: none; font-family: inherit; } .lin-sel, .lin-q { background: #0c1220; border: 1px solid #243152; border-radius: 4px; font-size: 11px; color: #8b9bb4; padding: 6px 8px; outline: none; font-family: inherit; }
.lin-q { color: #d7e1f0; padding: 6px 10px; width: 220px; } .lin-q { color: #d7e1f0; padding: 6px 10px; width: 220px; }
+13 -5
View File
@@ -31,10 +31,18 @@ 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,))
failed = self.db.fetchall( # HEL-529 fix 2: "待处理"只保留当前仍未恢复的最新异常。一条失败/停滞
"SELECT * FROM batches WHERE trade_date = ? AND state IN ('failed','staged')", # 批次若已被同数据集更晚的成功批次或当日发布解决,就不再是当前故障,
(today,), # 历史记录仍完整保留在 batches/audit 明细里,不在此重复展示。
) 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",
) )
@@ -45,7 +53,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": failed, "anomalies": anomalies,
"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,6 +38,7 @@ 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,
@@ -46,6 +47,7 @@ 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,
@@ -54,6 +56,7 @@ 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,
@@ -62,6 +65,7 @@ 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,
@@ -70,6 +74,7 @@ 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,
@@ -78,6 +83,7 @@ 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,
@@ -86,6 +92,7 @@ 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,
@@ -94,6 +101,7 @@ 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,
@@ -102,6 +110,7 @@ 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,
@@ -110,6 +119,7 @@ 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,
@@ -118,6 +128,7 @@ 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,
@@ -126,6 +137,7 @@ 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",
@@ -134,6 +146,7 @@ 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",
@@ -142,6 +155,7 @@ 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,
@@ -150,6 +164,7 @@ 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,
@@ -158,6 +173,7 @@ 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,
@@ -166,6 +182,7 @@ 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,6 +99,58 @@ 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 "
@@ -1720,6 +1772,7 @@ 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 = ?",
+60 -29
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"]}, {"interface": "trade_cal", "capability": "交易日历", "datasets": ["calendar"], "group": "日历 / 主档"},
{"interface": "stock_basic", "capability": "股票主档", "datasets": ["stocks"]}, {"interface": "stock_basic", "capability": "股票主档", "datasets": ["stocks"], "group": "日历 / 主档"},
{"interface": "daily", "capability": "个股日K", "datasets": ["daily"]}, {"interface": "daily", "capability": "个股日K", "datasets": ["daily"], "group": "盘后 A 批"},
{"interface": "adj_factor", "capability": "复权因子", "datasets": ["daily"]}, {"interface": "adj_factor", "capability": "复权因子", "datasets": ["daily"], "group": "盘后 A 批"},
{"interface": "daily_basic", "capability": "估值", "datasets": ["valuation"]}, {"interface": "daily_basic", "capability": "估值", "datasets": ["valuation"], "group": "盘后 A 批"},
{"interface": "index_daily", "capability": "指数日K", "datasets": ["index_daily"]}, {"interface": "index_daily", "capability": "指数日K", "datasets": ["index_daily"], "group": "指数 B 批"},
{"interface": "moneyflow", "capability": "资金流", "datasets": ["moneyflow"]}, {"interface": "moneyflow", "capability": "资金流", "datasets": ["moneyflow"], "group": "盘后 A 批"},
{"interface": "stk_auction", "capability": "集合竞价", "datasets": ["auction"]}, {"interface": "stk_auction", "capability": "集合竞价", "datasets": ["auction"], "group": "盘后 A 批"},
{"interface": "limit_list_d", "capability": "涨跌停池", "datasets": ["limit_events"]}, {"interface": "limit_list_d", "capability": "涨跌停池", "datasets": ["limit_events"], "group": "扩展软批"},
{"interface": "ths_hot", "capability": "同花顺人气榜", "datasets": ["popularity"]}, {"interface": "ths_hot", "capability": "同花顺人气榜", "datasets": ["popularity"], "group": "扩展软批"},
{"interface": "dc_hot", "capability": "东方财富人气榜", "datasets": ["popularity"]}, {"interface": "dc_hot", "capability": "东方财富人气榜", "datasets": ["popularity"], "group": "扩展软批"},
{"interface": "hm_detail", "capability": "龙虎榜游资明细", "datasets": ["dragon_tiger"]}, {"interface": "hm_detail", "capability": "龙虎榜游资明细", "datasets": ["dragon_tiger"], "group": "扩展软批"},
{"interface": "ths_daily", "capability": "同花顺概念行情", "datasets": ["sector_daily"]}, {"interface": "ths_daily", "capability": "同花顺概念行情", "datasets": ["sector_daily"], "group": "扩展软批"},
{"interface": "dc_index", "capability": "东方财富板块行情", "datasets": ["sector_daily"]}, {"interface": "dc_index", "capability": "东方财富板块行情", "datasets": ["sector_daily"], "group": "扩展软批"},
{"interface": "sw_daily", "capability": "申万行业行情", "datasets": ["sector_daily"]}, {"interface": "sw_daily", "capability": "申万行业行情", "datasets": ["sector_daily"], "group": "扩展软批"},
], ],
}, },
{ {
@@ -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"]}, {"interface": "indices", "capability": "指数实时报价", "datasets": ["index_quotes"], "group": "实时快照"},
{"interface": "market_quotes", "capability": "全市场实时快照", "datasets": ["quotes_latest"]}, {"interface": "market_quotes", "capability": "全市场实时快照", "datasets": ["quotes_latest"], "group": "实时快照"},
{"interface": "named_quotes", "capability": "指定个股实时报价", "datasets": ["quotes_latest"]}, {"interface": "named_quotes", "capability": "指定个股实时报价", "datasets": ["quotes_latest"], "group": "实时快照"},
{"interface": "sector_quote", "capability": "申万板块实时报价(单个)", "datasets": ["sectors_quote"]}, {"interface": "sector_quote", "capability": "申万板块实时报价(单个)", "datasets": ["sectors_quote"], "group": "实时快照"},
{"interface": "sector_quotes_batch", "capability": "申万板块批量报价(预热)", "datasets": ["sectors_quote"]}, {"interface": "sector_quotes_batch", "capability": "申万板块批量报价(预热)", "datasets": ["sectors_quote"], "group": "实时快照"},
{"interface": "limit_pool", "capability": "涨停/炸板池(盘中)", "datasets": ["limit_pool"]}, {"interface": "limit_pool", "capability": "涨停/炸板池(盘中)", "datasets": ["limit_pool"], "group": "实时快照"},
{"interface": "intraday", "capability": "分时走势", "datasets": ["intraday_points"]}, {"interface": "intraday", "capability": "分时走势", "datasets": ["intraday_points"], "group": "实时快照"},
], ],
}, },
{ {
@@ -71,13 +71,14 @@ 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"]}, {"interface": "indices", "capability": "指数实时报价(东财失败时备用)", "datasets": ["index_quotes"], "group": "实时备援"},
{ {
"interface": "market_quotes_fallback", "interface": "market_quotes_fallback",
"capability": "全市场快照(备用;按本地股票主档逐只请求拼接)", "capability": "全市场快照(备用;按本地股票主档逐只请求拼接)",
"datasets": ["quotes_latest"], "datasets": ["quotes_latest"],
"group": "实时备援",
}, },
{"interface": "named_quotes", "capability": "指定个股实时报价(东财失败时备用)", "datasets": ["quotes_latest"]}, {"interface": "named_quotes", "capability": "指定个股实时报价(东财失败时备用)", "datasets": ["quotes_latest"], "group": "实时备援"},
], ],
}, },
{ {
@@ -87,11 +88,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"]}, {"interface": "wencai", "capability": "问财自然语言选股", "datasets": ["ifind_wencai"], "group": "预留接口"},
{"interface": "snapshots", "capability": "快照", "datasets": ["ifind_snapshots"]}, {"interface": "snapshots", "capability": "快照", "datasets": ["ifind_snapshots"], "group": "预留接口"},
{"interface": "history", "capability": "历史行情", "datasets": ["ifind_history"]}, {"interface": "history", "capability": "历史行情", "datasets": ["ifind_history"], "group": "预留接口"},
{"interface": "realtime", "capability": "实时行情", "datasets": ["ifind_realtime"]}, {"interface": "realtime", "capability": "实时行情", "datasets": ["ifind_realtime"], "group": "预留接口"},
{"interface": "intraday", "capability": "分时(高频)", "datasets": ["ifind_intraday"]}, {"interface": "intraday", "capability": "分时(高频)", "datasets": ["ifind_intraday"], "group": "预留接口"},
], ],
}, },
{ {
@@ -161,5 +162,35 @@ 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
@@ -0,0 +1,213 @@
"""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()