diff --git a/backend/data/gateway.py b/backend/data/gateway.py index 4abdac5..6926585 100644 --- a/backend/data/gateway.py +++ b/backend/data/gateway.py @@ -37,7 +37,9 @@ class DataGateway: ) -> TushareClient: if dataset_id: self.policy.assert_allowed(dataset_id, "tushare", usage) - return DatahubAwareTushareClient(self.tushare_provider.client(), self.datahub) + legacy = self.tushare_provider.client() + legacy.realtime_aggregator = self.realtime_observer + return DatahubAwareTushareClient(legacy, self.datahub) def dataset_status(self, trade_date: str) -> list[dict[str, Any]] | None: return self.datahub.dataset_status(trade_date) diff --git a/backend/data/providers/tushare_dashboard.py b/backend/data/providers/tushare_dashboard.py index 2b1404a..d0b0d36 100644 --- a/backend/data/providers/tushare_dashboard.py +++ b/backend/data/providers/tushare_dashboard.py @@ -128,7 +128,7 @@ class DashboardMixin: ) if not codes: raise TushareError("No active stock codes available for rt_k") - quotes = self.query("rt_k", {"ts_code": codes}) + quotes, quote_source = self._load_realtime_quotes(codes, trade_date) if not quotes: raise TushareError(f"No realtime data returned for {trade_date}") @@ -186,12 +186,28 @@ class DashboardMixin: previous_sectors = _build_sectors(previous_limits) now = self._now() market_status = _realtime_market_status(now.time().replace(tzinfo=None)) + if quote_source == "eastmoney_clist": + notice = ( + "盘中行情由东财免费实时快照计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。" + ) + source_name = "eastmoney" + elif quote_source == "tencent_qt": + notice = ( + "盘中行情由腾讯免费实时行情计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。" + ) + source_name = "tencent" + else: + notice = ( + "盘中行情由 Tushare rt_k 实时计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。" + ) + source_name = "tushare" dashboard = { "meta": { "requested_date": _display_date(requested_date), "trade_date": _display_date(trade_date), "previous_trade_date": _display_date(previous_trade_date), - "source": "tushare", + "source": source_name, + "quote_source": quote_source, "mode": "realtime", "realtime": True, "market_status": market_status, @@ -199,7 +215,8 @@ class DashboardMixin: "auto_refresh": False, "quote_count": len(daily), "updated_at": now.isoformat(timespec="seconds"), - "notice": "盘中行情由 Tushare rt_k 实时计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。", + "notice": notice, + "indices": self._free_realtime_indices() if quote_source != "tushare_rt_k" else [], }, "overview": _build_overview(daily, up_rows, down_rows, broken_rows), "limits": limits, @@ -213,6 +230,67 @@ class DashboardMixin: } return apply_sentiment_to_dashboard(dashboard) + def _realtime_aggregator(self): + aggregator = getattr(self, "realtime_aggregator", None) + if aggregator is None: + raise TushareError("免费实时源未配置") + return aggregator + + def _load_realtime_quotes( + self, + codes: str, + trade_date: str, + ) -> tuple[list[dict[str, Any]], str]: + rt_error = "" + try: + quotes = self.query("rt_k", {"ts_code": codes}) + if quotes: + return list(quotes), "tushare_rt_k" + rt_error = f"No realtime data returned for {trade_date}" + except TushareError as exc: + rt_error = str(exc) + try: + quotes, quote_source = self._free_realtime_quotes(trade_date, codes) + except Exception as exc: + raise TushareError( + f"当天盘中实时行情不可用:rt_k={rt_error};免费源={exc}" + ) from exc + if not quotes: + raise TushareError( + f"当天盘中实时行情不可用:rt_k={rt_error};免费源=empty" + ) + return quotes, quote_source + + def _free_realtime_quotes( + self, + trade_date: str, + codes: str = "", + ) -> tuple[list[dict[str, Any]], str]: + aggregator = self._realtime_aggregator() + last_error = "" + try: + quotes = aggregator.eastmoney_market_quotes(expected_date=trade_date) + if quotes: + return quotes, "eastmoney_clist" + except Exception as exc: + last_error = str(exc) + code_list = [item for item in str(codes or "").split(",") if item] + try: + quotes = aggregator.tencent_market_quotes(code_list, expected_date=trade_date) + except Exception as exc: + raise TushareError( + f"eastmoney={last_error or 'empty'};tencent={exc}" + ) from exc + if not quotes: + raise TushareError(f"eastmoney={last_error or 'empty'};tencent=empty") + return quotes, "tencent_qt" + + def _free_realtime_indices(self) -> list[dict[str, Any]]: + try: + return self._realtime_aggregator().eastmoney_indices() + except Exception: + return [] + def _load_realtime_reference( self, trade_date: str, diff --git a/backend/data/providers/tushare_indices.py b/backend/data/providers/tushare_indices.py index 7087628..ce1a223 100644 --- a/backend/data/providers/tushare_indices.py +++ b/backend/data/providers/tushare_indices.py @@ -59,6 +59,12 @@ class IndexMixin: } def realtime_market_indices(self, requested_date: str) -> dict[str, Any]: + try: + return self._tushare_realtime_market_indices(requested_date) + except TushareError: + return self._free_realtime_market_indices(requested_date) + + def _tushare_realtime_market_indices(self, requested_date: str) -> dict[str, Any]: trade_date, _ = self.resolve_trade_context(requested_date) index_names = { "000001.SH": "上证指数", @@ -116,3 +122,52 @@ class IndexMixin: "average_return_20d": 0, }, } + + def _free_realtime_market_indices(self, requested_date: str) -> dict[str, Any]: + trade_date, _ = self.resolve_trade_context(requested_date) + aggregator = getattr(self, "realtime_aggregator", None) + if aggregator is None: + raise TushareError("免费实时源未配置") + quotes = aggregator.eastmoney_indices() + index_names = { + "000001": ("000001.SH", "上证指数"), + "399001": ("399001.SZ", "深证成指"), + "399006": ("399006.SZ", "创业板指"), + } + indices = [] + for quote in quotes: + mapped = index_names.get(str(quote.get("code") or "")) + if not mapped: + continue + ts_code, name = mapped + close = _number(quote.get("price")) + previous_close = _number(quote.get("previous_close")) + if close <= 0 or previous_close <= 0: + continue + indices.append( + { + "ts_code": ts_code, + "name": str(quote.get("name") or name).strip(), + "trade_date": trade_date, + "close": close, + "pct_chg": round(_number(quote.get("change")) or (close / previous_close - 1) * 100, 3), + "return_5d": 0, + "amount_billion": round(_number(quote.get("amount_billion")), 2), + "quote_time": quote.get("quote_time") or "", + "source": quote.get("source") or "eastmoney_push2", + } + ) + if len(indices) != 3: + raise TushareError("Realtime index quotes are incomplete") + return { + "trade_date": trade_date, + "source": "eastmoney_push2", + "realtime": True, + "precise": True, + "indices": indices, + "aggregate": { + "average_pct_chg": round(sum(item["pct_chg"] for item in indices) / len(indices), 3), + "average_return_5d": 0, + "average_return_20d": 0, + }, + } diff --git a/backend/data/realtime.py b/backend/data/realtime.py index d566df2..7222571 100644 --- a/backend/data/realtime.py +++ b/backend/data/realtime.py @@ -20,7 +20,17 @@ class RealtimeAggregateError(RuntimeError): EASTMONEY_INDEX_URL = "https://push2.eastmoney.com/api/qt/ulist.np/get" EASTMONEY_SECTOR_URL = "https://push2.eastmoney.com/api/qt/clist/get" +EASTMONEY_A_SHARE_BOARDS = ( + "m:0+t:6", + "m:0+t:80", + "m:1+t:2", + "m:1+t:23", + "m:0+t:81", +) +EASTMONEY_QUOTE_FIELDS = "f12,f13,f14,f2,f3,f4,f5,f6,f15,f16,f17,f18,f8,f124" +EASTMONEY_MARKET_PAGE_SIZE = 100 TENCENT_INDEX_URL = "https://qt.gtimg.cn/q=sh000001,sz399001,sz399006" +TENCENT_QUOTE_URL = "https://qt.gtimg.cn/q=" THS_LIMIT_URL = "https://data.10jqka.com.cn/dataapi/limit_up/limit_up_pool" XGB_POOL_URL = "https://flash-api.xuangubao.cn/api/pool/detail" BROWSER_USER_AGENT = ( @@ -134,6 +144,145 @@ class WebRealtimeAggregator: raise RealtimeAggregateError(f"Eastmoney returned {len(result)}/3 indices") return result + def eastmoney_market_quotes(self, expected_date: str = "") -> list[dict[str, Any]]: + """Full A-share snapshot via Eastmoney clist, used when Tushare rt_k is unavailable.""" + now = time.time() + cache_key = "assembled:eastmoney_market" + with self._response_cache_lock: + cached = self._response_cache.get(cache_key) + cache_age = now - float((cached or {}).get("created_at") or 0) + if cached and cache_age <= min(20, self.response_cache_ttl_seconds): + quotes = list(cached.get("payload") or []) + return self._filter_quotes_by_date(quotes, expected_date) + + rows: list[dict[str, Any]] = [] + board_errors: list[str] = [] + for board in EASTMONEY_A_SHARE_BOARDS: + try: + rows.extend(self._eastmoney_board_quotes(board)) + except Exception as exc: + board_errors.append(f"{board}:{exc}") + quotes = [] + seen: set[str] = set() + for row in rows: + quote = _normalize_eastmoney_quote(row) + ts_code = str((quote or {}).get("ts_code") or "") + if not quote or ts_code in seen: + continue + seen.add(ts_code) + quotes.append(quote) + if len(quotes) < 200: + detail = f";{'; '.join(board_errors)}" if board_errors else "" + raise RealtimeAggregateError( + f"Eastmoney market snapshot too small: {len(quotes)}{detail}" + ) + quotes = self._filter_quotes_by_date(quotes, expected_date) + with self._response_cache_lock: + self._response_cache[cache_key] = {"created_at": now, "payload": quotes} + return quotes + + def _eastmoney_board_quotes(self, board: str) -> list[dict[str, Any]]: + first = self._eastmoney_market_page(board, 1) + data = first.get("data") or {} + rows = _diff_rows(data) + total = int(_number(data.get("total"))) + page_count = 1 + if total > 0: + page_count = max(1, (total + EASTMONEY_MARKET_PAGE_SIZE - 1) // EASTMONEY_MARKET_PAGE_SIZE) + for page in range(2, min(page_count, 40) + 1): + payload = self._eastmoney_market_page(board, page) + rows.extend(_diff_rows(payload.get("data") or {})) + return rows + + def _eastmoney_market_page(self, board: str, page: int) -> dict[str, Any]: + return self._get_json( + EASTMONEY_SECTOR_URL, + { + "pn": str(page), + "pz": str(EASTMONEY_MARKET_PAGE_SIZE), + "po": "1", + "np": "1", + "fltt": "2", + "invt": "2", + "fid": "f12", + "fs": board, + "fields": EASTMONEY_QUOTE_FIELDS, + }, + referer="https://quote.eastmoney.com/center/gridlist.html", + ) + + def _filter_quotes_by_date( + self, + quotes: list[dict[str, Any]], + expected_date: str, + ) -> list[dict[str, Any]]: + want = str(expected_date or "").replace("-", "") + if not want or not quotes: + return quotes + dated = [item for item in quotes if str(item.get("quote_date") or "") == want] + if dated and len(dated) >= max(100, int(len(quotes) * 0.2)): + return dated + if dated: + return dated + if all(not item.get("quote_date") for item in quotes): + return quotes + raise RealtimeAggregateError(f"Eastmoney quotes are not for {want}") + + def tencent_market_quotes( + self, + codes: list[str], + expected_date: str = "", + ) -> list[dict[str, Any]]: + symbols: list[str] = [] + seen: set[str] = set() + for raw in codes: + ts = str(raw or "").strip().upper() + if not ts: + continue + symbol = ts.split(".")[0] + if not symbol.isdigit() or len(symbol) != 6 or symbol in seen: + continue + seen.add(symbol) + if ts.endswith(".SH") or symbol.startswith(("5", "6", "9")): + symbols.append(f"sh{symbol}") + elif ts.endswith(".BJ") or symbol.startswith(("4", "8")): + symbols.append(f"bj{symbol}") + else: + symbols.append(f"sz{symbol}") + if not symbols: + raise RealtimeAggregateError("No stock codes available for Tencent quotes") + + quotes: list[dict[str, Any]] = [] + batch_size = 80 + + def load_batch(batch: list[str]) -> list[dict[str, Any]]: + raw, _cache_age = self._get_text( + f"{TENCENT_QUOTE_URL}{','.join(batch)}", + referer="https://gu.qq.com/", + encoding="gb18030", + ) + return [ + quote + for line in raw.splitlines() + if (quote := _parse_tencent_stock_quote(line)) + ] + + batches = [symbols[index:index + batch_size] for index in range(0, len(symbols), batch_size)] + errors: list[str] = [] + with ThreadPoolExecutor(max_workers=4) as executor: + for result in executor.map(self._capture, [lambda batch=batch: load_batch(batch) for batch in batches]): + rows, status = result + if status.get("ok") and rows: + quotes.extend(rows) + elif not status.get("ok"): + errors.append(str(status.get("error") or "batch failed")) + if len(quotes) < 200: + detail = f";{'; '.join(errors[:3])}" if errors else "" + raise RealtimeAggregateError( + f"Tencent market snapshot too small: {len(quotes)}{detail}" + ) + return self._filter_quotes_by_date(quotes, expected_date) + def tencent_indices(self) -> list[dict[str, Any]]: raw, cache_age = self._get_text( TENCENT_INDEX_URL, @@ -397,6 +546,94 @@ class WebRealtimeAggregator: ) from last_error +def _diff_rows(data: dict[str, Any]) -> list[dict[str, Any]]: + diff = data.get("diff") or [] + if isinstance(diff, dict): + return [row for row in diff.values() if isinstance(row, dict)] + return [row for row in diff if isinstance(row, dict)] + + +def _parse_tencent_stock_quote(line: str) -> dict[str, Any] | None: + if '="' not in line: + return None + prefix, payload = line.split('="', 1) + fields = payload.rsplit('";', 1)[0].split("~") + if len(fields) < 38: + return None + symbol = fields[2] + if not symbol.isdigit() or len(symbol) != 6: + return None + close = _number(fields[3]) + previous_close = _number(fields[4]) + if close <= 0 or previous_close <= 0: + return None + marker = prefix.lower() + if "sh" in marker: + ts_code = f"{symbol}.SH" + elif "bj" in marker: + ts_code = f"{symbol}.BJ" + else: + ts_code = f"{symbol}.SZ" + try: + quote_time = datetime.strptime(fields[30], "%Y%m%d%H%M%S") + quote_date = quote_time.strftime("%Y%m%d") + epoch = int(quote_time.timestamp()) + except ValueError: + quote_date = "" + epoch = 0 + return { + "ts_code": ts_code, + "name": fields[1] or symbol, + "pre_close": previous_close, + "open": _number(fields[5]), + "high": _number(fields[33]), + "low": _number(fields[34]), + "close": close, + "vol": _number(fields[6]) * 100, + "amount": _number(fields[37]) * 10000, + "num": 0, + "quote_date": quote_date, + "quote_time_epoch": epoch, + "source": "tencent_qt", + } + + +def _normalize_eastmoney_quote(row: dict[str, Any]) -> dict[str, Any] | None: + symbol = str(row.get("f12") or "").strip() + if not symbol.isdigit() or len(symbol) != 6: + return None + close = _number(row.get("f2")) + previous_close = _number(row.get("f18")) + if close <= 0 or previous_close <= 0: + return None + market = int(_number(row.get("f13"))) + if market == 1 or symbol.startswith(("5", "6", "9")): + ts_code = f"{symbol}.SH" + elif symbol.startswith(("4", "8")): + ts_code = f"{symbol}.BJ" + else: + ts_code = f"{symbol}.SZ" + epoch = int(_number(row.get("f124"))) + quote_date = "" + if epoch > 0: + quote_date = datetime.fromtimestamp(epoch).astimezone().strftime("%Y%m%d") + return { + "ts_code": ts_code, + "name": row.get("f14") or symbol, + "pre_close": previous_close, + "open": _number(row.get("f17")), + "high": _number(row.get("f15")), + "low": _number(row.get("f16")), + "close": close, + "vol": _number(row.get("f5")) * 100, + "amount": _number(row.get("f6")), + "num": 0, + "quote_date": quote_date, + "quote_time_epoch": epoch, + "source": "eastmoney_clist", + } + + def _normalize_sector(value: Any) -> str: text = str(value or "").strip().replace(" ", "") for suffix in ("板块", "概念", "行业", "Ⅱ", "Ⅲ", "(A股)", "(A股)"): diff --git a/backend/features/market/service.py b/backend/features/market/service.py index 09756c2..07c9385 100644 --- a/backend/features/market/service.py +++ b/backend/features/market/service.py @@ -63,7 +63,11 @@ class MarketServiceMixin: if gateway is not None: return gateway.tushare() # Compatibility for isolated legacy unit-test service stubs. - return TushareClient(self.token) + client = TushareClient(self.token) + aggregator = getattr(self, "realtime_aggregator", None) + if aggregator is not None: + client.realtime_aggregator = aggregator + return client def _now(self) -> datetime: clock = getattr(self, "clock", None) @@ -289,7 +293,10 @@ class MarketServiceMixin: raise TushareError("公共行情尚未配置") dashboard = self._tushare_client().dashboard(normalized_date) meta = dashboard.setdefault("meta", {}) + quote_source = str(meta.get("quote_source") or "") meta["source"] = source + if quote_source: + meta["quote_source"] = quote_source meta["requested_date"] = self._display_compact_date(normalized_date) if meta.get("limit_data_source") == "derived": meta.setdefault( diff --git a/config/architecture-inventory.json b/config/architecture-inventory.json index e05e616..6f20e27 100644 --- a/config/architecture-inventory.json +++ b/config/architecture-inventory.json @@ -222,12 +222,12 @@ { "provider": "eastmoney", "path": "backend/data/realtime.py", - "runtime_role": "isolated realtime observation" + "runtime_role": "isolated realtime observation and intraday dashboard fallback" }, { "provider": "tencent", "path": "backend/data/realtime.py", - "runtime_role": "index observation fallback" + "runtime_role": "index observation and intraday quote fallback" } ], "provider_domains": [ @@ -508,8 +508,8 @@ }, { "path": "backend/data/providers/tushare_dashboard.py", - "bytes": 28327, - "lines": 654 + "bytes": 31361, + "lines": 732 }, { "path": "backend/data/providers/tushare_industries.py", @@ -631,6 +631,11 @@ "bytes": 8357, "lines": 116 }, + { + "path": "backend/data/providers/tushare_indices.py", + "bytes": 7823, + "lines": 173 + }, { "path": "backend/features/screener/formula.py", "bytes": 6983, @@ -686,11 +691,6 @@ "bytes": 5690, "lines": 124 }, - { - "path": "backend/data/providers/tushare_indices.py", - "bytes": 5451, - "lines": 118 - }, { "path": "frontend/pages.config.js", "bytes": 5385, diff --git a/docs/governance/architecture-inventory.json b/docs/governance/architecture-inventory.json index fd68534..c01fcc7 100644 --- a/docs/governance/architecture-inventory.json +++ b/docs/governance/architecture-inventory.json @@ -213,12 +213,12 @@ { "provider": "eastmoney", "path": "realtime_aggregator.py", - "runtime_role": "isolated realtime observation" + "runtime_role": "isolated realtime observation and intraday dashboard fallback" }, { "provider": "tencent", "path": "realtime_aggregator.py", - "runtime_role": "index observation fallback" + "runtime_role": "index observation and intraday quote fallback" } ], "llm_entrypoints": [ diff --git a/tests/test_admin_refresh_status.py b/tests/test_admin_refresh_status.py index 3c96987..e64a25e 100644 --- a/tests/test_admin_refresh_status.py +++ b/tests/test_admin_refresh_status.py @@ -150,6 +150,32 @@ class FakeRealtimeTodayClient: return requested, "20260907" +class FakeFreeRealtimeTodayClient: + def dashboard(self, trade_date: str): + return { + "meta": { + "trade_date": f"{trade_date[:4]}-{trade_date[4:6]}-{trade_date[6:8]}", + "requested_date": f"{trade_date[:4]}-{trade_date[4:6]}-{trade_date[6:8]}", + "realtime": True, + "mode": "realtime", + "quote_source": "eastmoney_clist", + "source": "eastmoney", + "market_status": "trading", + "notice": "盘中行情由东财免费实时快照计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。", + "updated_at": datetime.now().astimezone().isoformat(timespec="seconds"), + "indices": [{"code": "000001", "price": 3800.1, "change": 0.5}], + }, + "overview": {"limit_up_count": 18, "up_count": 2100, "amount_billion": 12345.6}, + "limits": [{"code": "000001"}], + "broken": [], + "down_limits": [], + "yesterday_limits": [], + } + + def resolve_trade_context(self, requested: str): + return requested, "20260907" + + class SyncHarness(MarketServiceMixin): def __init__(self, client, latest=None, clock=None): self.configured = True @@ -204,6 +230,28 @@ class DashboardFreshnessTests(unittest.TestCase): self.assertNotIn("今日数据正在准备", meta.get("display_notice") or "") self.assertEqual(harness.database.saved[0][0], today) + def test_intraday_free_source_keeps_today_and_indices(self): + today = TRADE_DAY.strftime("%Y%m%d") + latest = { + "meta": {"trade_date": "2026-09-07", "source": "tushare"}, + "overview": {"limit_up_count": 20}, + } + harness = SyncHarness( + FakeFreeRealtimeTodayClient(), + latest, + clock=lambda: at_clock(10, 5), + ) + payload = harness.sync_dashboard(today) + meta = payload["meta"] + self.assertFalse(meta.get("carried_forward")) + self.assertTrue(meta["realtime"]) + self.assertEqual(meta["data_status"], "intraday") + self.assertEqual(str(meta["trade_date"]).replace("-", ""), today) + self.assertEqual(meta["quote_source"], "eastmoney_clist") + self.assertEqual(payload["overview"]["amount_billion"], 12345.6) + self.assertEqual(meta["indices"][0]["price"], 3800.1) + self.assertEqual(harness.database.saved[0][0], today) + def test_intraday_missing_quotes_do_not_carry_yesterday(self): today = TRADE_DAY.strftime("%Y%m%d") latest = { diff --git a/tests/test_heaven_realtime.py b/tests/test_heaven_realtime.py index 3b9752f..323eeeb 100644 --- a/tests/test_heaven_realtime.py +++ b/tests/test_heaven_realtime.py @@ -3,6 +3,7 @@ from __future__ import annotations import http.client import json import unittest +from datetime import datetime from unittest.mock import MagicMock, patch from backend.data.realtime import WebRealtimeAggregator @@ -377,6 +378,43 @@ class RealtimeAggregatorTests(unittest.TestCase): self.assertEqual(rows[0]["quote_time"][:10], "2026-07-20") self.assertAlmostEqual(rows[0]["amount_billion"], 12946.52) + @patch.object(WebRealtimeAggregator, "_get_json") + def test_eastmoney_market_quotes_normalize_and_keep_expected_date(self, get_json: MagicMock): + epoch = datetime(2026, 7, 20, 10, 5).timestamp() + rows = [] + for index in range(200): + sz = index < 100 + rows.append( + { + "f12": f"{index:06d}" if sz else f"{600000 + index - 100:06d}", + "f13": 0 if sz else 1, + "f14": f"股票{index}", + "f2": 11.2, + "f3": 2.0, + "f5": 10, + "f6": 50000000, + "f15": 11.3, + "f16": 11.0, + "f17": 11.1, + "f18": 11.0, + "f124": epoch, + } + ) + def fake_get_json(_url, params, referer=""): + page = int(params.get("pn") or 1) + start = (page - 1) * 100 + return {"rc": 0, "data": {"total": 200, "diff": rows[start:start + 100]}} + + get_json.side_effect = fake_get_json + aggregator = WebRealtimeAggregator() + aggregator._response_cache.clear() + quotes = aggregator.eastmoney_market_quotes("20260720") + self.assertEqual(len(quotes), 200) + self.assertEqual(quotes[0]["ts_code"], "000000.SZ") + self.assertTrue(quotes[100]["ts_code"].endswith(".SH")) + self.assertEqual(quotes[0]["vol"], 1000) + self.assertEqual(quotes[0]["quote_date"], "20260720") + if __name__ == "__main__": unittest.main() diff --git a/tests/test_realtime_dashboard.py b/tests/test_realtime_dashboard.py index 8943df6..f5c32a3 100644 --- a/tests/test_realtime_dashboard.py +++ b/tests/test_realtime_dashboard.py @@ -5,6 +5,12 @@ from datetime import datetime, timedelta, timezone from backend.data.providers.tushare_client import TushareClient from backend.data.providers.tushare_helpers import calendar_is_open +from backend.data.providers.tushare_transport import TushareError +from backend.data.realtime import ( + RealtimeAggregateError, + _normalize_eastmoney_quote, + _parse_tencent_stock_quote, +) class FakeRealtimeClient(TushareClient): @@ -83,6 +89,65 @@ class FakeRealtimeClient(TushareClient): raise AssertionError(f"Unexpected API call: {api_name} {params}") +FREE_QUOTES = [ + { + "ts_code": "000001.SZ", "name": "甲", "pre_close": 10.0, + "open": 10.1, "high": 11.0, "low": 10.0, "close": 11.0, + "vol": 1000, "amount": 100000000, "num": 10, + "quote_date": "20260720", + }, + { + "ts_code": "000002.SZ", "name": "乙", "pre_close": 20.0, + "open": 19.5, "high": 20.0, "low": 18.0, "close": 18.0, + "vol": 2000, "amount": 200000000, "num": 20, + "quote_date": "20260720", + }, + { + "ts_code": "000003.SZ", "name": "丙", "pre_close": 30.0, + "open": 31.0, "high": 33.0, "low": 30.0, "close": 32.0, + "vol": 3000, "amount": 300000000, "num": 30, + "quote_date": "20260720", + }, +] + + +class FakeFreeAggregator: + def __init__(self, quotes=None, fail=False): + self.quotes = list(quotes if quotes is not None else FREE_QUOTES) + self.fail = fail + self.calls = 0 + + def eastmoney_market_quotes(self, expected_date=""): + self.calls += 1 + if self.fail: + raise RealtimeAggregateError("eastmoney down") + if expected_date and self.quotes: + dated = [ + row for row in self.quotes + if str(row.get("quote_date") or "") == str(expected_date).replace("-", "") + ] + if dated: + return dated + return list(self.quotes) + + def tencent_market_quotes(self, codes, expected_date=""): + return self.eastmoney_market_quotes(expected_date) + + def eastmoney_indices(self): + return [ + { + "code": "000001", + "name": "上证指数", + "price": 3800.12, + "change": 0.85, + "previous_close": 3768.0, + "amount_billion": 4200.5, + "quote_time": "2026-07-20T10:05:00+08:00", + "source": "eastmoney_push2", + } + ] + + class RealtimeDashboardTests(unittest.TestCase): def setUp(self): TushareClient._realtime_reference_cache.clear() @@ -184,6 +249,121 @@ class RealtimeDashboardTests(unittest.TestCase): self.assertEqual(dashboard["meta"]["quote_count"], 3) self.assertEqual(dashboard["overview"]["limit_up_count"], 0) + def test_rt_k_permission_error_falls_back_to_free_quotes(self): + original_query = self.client.query + + def query(api_name, params=None, fields=""): + if api_name == "rt_k": + raise TushareError("没有接口访问权限") + return original_query(api_name, params, fields) + + self.client.query = query + self.client.realtime_aggregator = FakeFreeAggregator() + TushareClient._realtime_reference_cache.clear() + dashboard = self.client._realtime_dashboard("20260720", "20260720", "20260717") + + self.assertTrue(dashboard["meta"]["realtime"]) + self.assertEqual(dashboard["meta"]["quote_source"], "eastmoney_clist") + self.assertEqual(dashboard["meta"]["trade_date"], "2026-07-20") + self.assertEqual(dashboard["meta"]["quote_count"], 3) + self.assertEqual(dashboard["overview"]["limit_up_count"], 1) + self.assertEqual(dashboard["overview"]["limit_down_count"], 1) + self.assertEqual(dashboard["overview"]["amount_billion"], 6.0) + self.assertIn("东财免费实时", dashboard["meta"]["notice"]) + self.assertEqual(dashboard["meta"]["indices"][0]["price"], 3800.12) + + def test_rt_k_empty_result_falls_back_to_free_quotes(self): + original_query = self.client.query + + def query(api_name, params=None, fields=""): + if api_name == "rt_k": + return [] + return original_query(api_name, params, fields) + + self.client.query = query + self.client.realtime_aggregator = FakeFreeAggregator() + TushareClient._realtime_reference_cache.clear() + dashboard = self.client._realtime_dashboard("20260720", "20260720", "20260717") + self.assertEqual(dashboard["meta"]["quote_source"], "eastmoney_clist") + self.assertEqual(str(dashboard["meta"]["trade_date"]).replace("-", ""), "20260720") + + def test_rt_k_and_free_source_failure_keeps_today_error(self): + original_query = self.client.query + + def query(api_name, params=None, fields=""): + if api_name == "rt_k": + raise TushareError("没有接口访问权限") + return original_query(api_name, params, fields) + + self.client.query = query + self.client.realtime_aggregator = FakeFreeAggregator(fail=True) + TushareClient._realtime_reference_cache.clear() + with self.assertRaises(TushareError) as ctx: + self.client._realtime_dashboard("20260720", "20260720", "20260717") + self.assertIn("当天盘中实时行情不可用", str(ctx.exception)) + self.assertIn("没有接口访问权限", str(ctx.exception)) + + def test_rt_k_and_eastmoney_failure_falls_back_to_tencent(self): + original_query = self.client.query + + def query(api_name, params=None, fields=""): + if api_name == "rt_k": + raise TushareError("没有接口访问权限") + return original_query(api_name, params, fields) + + class TencentOnlyAggregator(FakeFreeAggregator): + def eastmoney_market_quotes(self, expected_date=""): + raise RealtimeAggregateError("eastmoney blocked") + + def tencent_market_quotes(self, codes, expected_date=""): + return list(FREE_QUOTES) + + self.client.query = query + self.client.realtime_aggregator = TencentOnlyAggregator() + TushareClient._realtime_reference_cache.clear() + dashboard = self.client._realtime_dashboard("20260720", "20260720", "20260717") + self.assertEqual(dashboard["meta"]["quote_source"], "tencent_qt") + self.assertEqual(str(dashboard["meta"]["trade_date"]).replace("-", ""), "20260720") + self.assertIn("腾讯免费实时", dashboard["meta"]["notice"]) + self.assertEqual(dashboard["overview"]["amount_billion"], 6.0) + + def test_normalize_eastmoney_quote_maps_units_and_exchange(self): + quote = _normalize_eastmoney_quote( + { + "f12": "600000", + "f13": 1, + "f14": "浦发银行", + "f2": 10.5, + "f5": 12.0, + "f6": 200000000, + "f15": 10.8, + "f16": 10.2, + "f17": 10.3, + "f18": 10.0, + "f124": 1752986700, + } + ) + self.assertEqual(quote["ts_code"], "600000.SH") + self.assertEqual(quote["vol"], 1200) + self.assertEqual(quote["close"], 10.5) + self.assertEqual(quote["pre_close"], 10.0) + self.assertEqual(quote["source"], "eastmoney_clist") + + def test_parse_tencent_stock_quote_keeps_today_and_units(self): + line = ( + 'v_sz000001="51~平安银行~000001~11.73~11.70~11.66~346232~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~0~' + '~20260720100500~0.03~0.26~11.79~11.65~11.73/346232/406045563~346232~40605~0.18~5.24~~11.79~11.65~1.20~' + '2276.29~2276.31~0.49~12.87~10.53~0.95~-3076~11.73~4.43~5.34~~~0.18~40604.5563~0.0000~0~";' + ) + quote = _parse_tencent_stock_quote(line) + self.assertEqual(quote["ts_code"], "000001.SZ") + self.assertEqual(quote["quote_date"], "20260720") + self.assertEqual(quote["close"], 11.73) + self.assertEqual(quote["pre_close"], 11.70) + self.assertEqual(quote["vol"], 34623200) + self.assertEqual(quote["amount"], 406050000) + self.assertEqual(quote["source"], "tencent_qt") + if __name__ == "__main__": unittest.main() diff --git a/tools/build_architecture_inventory.py b/tools/build_architecture_inventory.py index ddb12a3..9f8d222 100644 --- a/tools/build_architecture_inventory.py +++ b/tools/build_architecture_inventory.py @@ -221,8 +221,8 @@ def build() -> dict[str, Any]: {"provider": "datahub", "path": "backend/data/datahub/client.py", "runtime_role": "optional official EOD read path behind per-dataset flags"}, {"provider": "ifind", "path": "backend/data/providers/ifind_client.py", "runtime_role": "realtime, charts, snapshots, enrichment"}, {"provider": "eastmoney", "path": "backend/features/market/charts.py", "runtime_role": "display chart fallback"}, - {"provider": "eastmoney", "path": "backend/data/realtime.py", "runtime_role": "isolated realtime observation"}, - {"provider": "tencent", "path": "backend/data/realtime.py", "runtime_role": "index observation fallback"}, + {"provider": "eastmoney", "path": "backend/data/realtime.py", "runtime_role": "isolated realtime observation and intraday dashboard fallback"}, + {"provider": "tencent", "path": "backend/data/realtime.py", "runtime_role": "index observation and intraday quote fallback"}, ], "provider_domains": [ {"provider": "tushare", "path": "backend/data/providers/tushare_transport.py", "responsibility": "HTTP transport and provider errors"},