Files
xiaobai-review/xiaobai-datahub/tests/test_hel529_rework.py
T
施工员2号andmultica-agent 35ee43ea02 fix(HEL-529): 布局鲁棒、问师血缘修正与动效语义增强
布局:表格长错误/备注统一单行大白话摘要+截断(.ellip),完整原文放
title 悬停;当前来源列 nowrap 防逐字竖排;压力态(长错误/多接口/多观测行)
1920×920、1920×1080、1024 均无横向溢出、血缘 1080P 一屏。
血缘:按真实调用代码修正——ifind_wencai 的真实消费者是股票池事件补充
(pools/service.py:82),不是问师;新增 ifind_history 条目对应问师·趋势/
宏观思维模型的可选指数/ETF 矩阵 (mentor/service.py:433 ifind.history,
未配置时 fail-open);新增 3 项自动测试锁定映射并以源码调用点为证。
动效:延迟/成功率仅在真实采样值变化时翻动高亮(data-rate flash);
spark-end/血缘心跳按行/节点错峰相位;页面隐藏时动态循环不做任何工作;
LINK 呼吸(gnode-pulse)与真实调用脉冲(node-ping)语义区分保持。
测试:全套 194 项通过。

Co-authored-by: multica-agent <github@multica.ai>
2026-09-15 15:11:32 +08:00

261 lines
14 KiB
Python

"""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']}")
class LineageMentorIfindTests(_Base):
"""问师 → iFinD 血缘修正(第二轮返工):不按名称猜关系,按真实调用代码。"""
REPO = Path(__file__).resolve().parents[2]
def test_no_mentor_dependency_on_ifind_wencai(self) -> None:
items = self.hub.admin.lineage("20240902")["items"]
wencai = [i for i in items if i["dataset"] == "ifind_wencai"]
self.assertEqual(len(wencai), 1)
consumers = wencai[0]["known_consumers"]
for consumer in consumers:
self.assertNotIn("问师", consumer, f"wencai consumer must not be 问师: {consumer}")
all_consumers = " ".join(c for i in items for c in i["known_consumers"])
self.assertNotIn("问师(自然语言选股", all_consumers)
def test_wencai_real_consumer_is_pools_enrichment_with_code_evidence(self) -> None:
items = self.hub.admin.lineage("20240902")["items"]
wencai = [i for i in items if i["dataset"] == "ifind_wencai"][0]
self.assertTrue(any("股票池" in c for c in wencai["known_consumers"]))
# Real call evidence in the main-site source tree:
pools_src = (self.REPO / "backend" / "features" / "pools" / "service.py").read_text(encoding="utf-8")
self.assertIn("ifind.wencai(", pools_src)
self.assertIn("ifind_event_enrichment_v1", pools_src)
# And 问师 itself never calls wencai:
mentor_src = (self.REPO / "backend" / "features" / "mentor" / "service.py").read_text(encoding="utf-8")
self.assertNotIn(".wencai(", mentor_src)
def test_mentor_optional_ifind_history_subcapabilities(self) -> None:
items = self.hub.admin.lineage("20240902")["items"]
history = [i for i in items if i["dataset"] == "ifind_history"]
self.assertEqual(len(history), 1)
consumers = history[0]["known_consumers"]
self.assertTrue(any("趋势思维模型" in c for c in consumers))
self.assertTrue(any("宏观思维模型" in c for c in consumers))
# every consumer must be a 问师 sub-capability, not the whole board
for consumer in consumers:
self.assertIn("·", consumer, f"not a sub-capability mapping: {consumer}")
# Real call evidence: mentor builds market matrices via ifind.history
mentor_src = (self.REPO / "backend" / "features" / "mentor" / "service.py").read_text(encoding="utf-8")
self.assertIn("ifind.history(", mentor_src)
self.assertIn("_mentor_market_matrix", mentor_src)
self.assertIn("MENTOR_INDEX_UNIVERSE", mentor_src)
self.assertIn("MENTOR_ETF_UNIVERSE", mentor_src)
# Optional dependency: fails open when ifind is not configured
self.assertIn("if not ifind or not ifind.configured", mentor_src)
if __name__ == "__main__":
unittest.main()