"""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()