From 1c2f2ac057d74a3ae96d4fae74765c04a2381c18 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=80=BB=E5=B7=A5?= Date: Tue, 8 Sep 2026 12:03:30 +0800 Subject: [PATCH] =?UTF-8?q?feat(HEL-490):=20=E5=89=A9=E4=BD=99=E8=A1=8C?= =?UTF-8?q?=E6=83=85=E6=94=B9=E7=94=B1=E6=95=B0=E6=8D=AE=E4=B8=AD=E6=9E=A2?= =?UTF-8?q?=E4=B8=BB=E7=BA=BF=E8=B7=AF=E6=8F=90=E4=BE=9B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 正式页面以 8766 为主线路,旧接口只作故障备用;compose 钉死全部 DATAHUB_READ_*,避免现网残留 0 造成假完成。 Co-authored-by: Cursor Co-authored-by: multica-agent --- .env.example | 9 +- backend/data/datahub/bridge.py | 203 +++++++++++++++++- backend/data/datahub/route_state.py | 57 +++++ backend/data/gateway.py | 25 +++ backend/data/providers/tushare_dashboard.py | 56 ++++- backend/data/providers/tushare_indices.py | 79 ++++++- backend/features/market/charts.py | 112 ++++++++++ backend/features/system/service.py | 17 ++ compose.yaml | 16 ++ config/README.md | 9 +- config/architecture-inventory.json | 40 ++-- config/datahub.config.json | 30 +-- frontend/index.html | 1 + frontend/m/js/pages.js | 10 + frontend/shared/admin.js | 21 ++ tests/test_chart_data_provider.py | 32 ++- tests/test_datahub_bridge.py | 111 +++++++++- tests/test_http_dispatch.py | 2 +- tests/test_realtime_dashboard.py | 17 ++ xiaobai-datahub/README.md | 2 +- xiaobai-datahub/datahub/adapters/eastmoney.py | 106 ++++++++- xiaobai-datahub/datahub/realtime_serve.py | 29 ++- xiaobai-datahub/datahub/serving.py | 6 +- .../tests/test_realtime_intraday.py | 56 +++++ 24 files changed, 979 insertions(+), 67 deletions(-) create mode 100644 backend/data/datahub/route_state.py diff --git a/.env.example b/.env.example index 46c4c14..310686c 100644 --- a/.env.example +++ b/.env.example @@ -5,11 +5,10 @@ APP_ENCRYPTION_KEY= # the system settings; all accounts use the same backend market snapshot. TUSHARE_TOKEN=your_tushare_token_here -# Optional xiaobai-datahub client. All DATAHUB_READ_* / DATAHUB_SHADOW_* flags -# default off in config/datahub.config.json, so the website keeps using Tushare. -# Extended datasets (HEL-463): LIMIT_EVENTS POPULARITY DRAGON_TIGER SECTOR_DAILY -# QUOTES INDEX_QUOTES INTRADAY — plus first-batch CALENDAR STOCKS DAILY INDEX_DAILY -# VALUATION MONEYFLOW AUCTION STATUS. +# Official xiaobai-datahub client. Read flags default on in config/datahub.config.json. +# compose.yaml pins every DATAHUB_READ_* to 1 so leftover .env zeros cannot keep +# official pages on the old APIs. Old website APIs are emergency fallback only. +# DATAHUB_SHADOW_* can still override a single dataset. DATAHUB_BASE_URL=http://127.0.0.1:8766 DATAHUB_TOKEN= diff --git a/backend/data/datahub/bridge.py b/backend/data/datahub/bridge.py index df86306..81f5a20 100644 --- a/backend/data/datahub/bridge.py +++ b/backend/data/datahub/bridge.py @@ -16,6 +16,7 @@ from backend.data.datahub.native import ( yyyymmdd, ) from backend.data.datahub.redact import redact_text, redact_value +from backend.data.datahub.route_state import LEDGER from backend.data.datahub.settings import DatahubSettings from backend.data.providers.tushare_client import TushareClient @@ -125,6 +126,7 @@ class DatahubBridge: raise DatahubError("EMPTY", "datahub intraday empty") if (response.meta or {}).get("stale"): raise DatahubError("STALE", "datahub intraday stale") + self._record_route("intraday", "datahub", str((response.meta or {}).get("source") or "datahub")) return { "entity_type": str(data.get("entity_type") or "stock"), "identifier": str(data.get("identifier") or code), @@ -139,6 +141,106 @@ class DatahubBridge: self._log_failure("intraday", exc) return None + def try_market_quotes(self, trade_date: str = "") -> list[dict[str, Any]] | None: + return self._try_quote_rows("quotes", {}, expected_date=trade_date, minimum=200) + + def try_quotes(self, codes: list[str]) -> list[dict[str, Any]] | None: + cleaned = [str(item or "").strip() for item in codes if str(item or "").strip()] + if not cleaned: + return None + return self._try_quote_rows("quotes", {"codes": ",".join(cleaned[:60])}, minimum=1) + + def try_index_quotes(self) -> list[dict[str, Any]] | None: + flags = self.settings.flags("index_quotes") + if not flags.read: + return None + try: + response = self.client.index_quotes() + rows = [dict(item) for item in (response.data or []) if isinstance(item, dict)] + if len(rows) < 3: + raise DatahubError("EMPTY", "datahub index quotes incomplete") + if (response.meta or {}).get("stale"): + raise DatahubError("STALE", "datahub index quotes stale") + self._record_route( + "index_quotes", + "datahub", + str((response.meta or {}).get("source") or "datahub"), + ) + return rows + except Exception as exc: + self._log_failure("index_quotes", exc) + return None + + def try_daily_chart( + self, + code: str, + end_date: str, + limit: int = 90, + dataset: str = "daily", + ) -> list[dict[str, Any]] | None: + flags = self.settings.flags(dataset) + if not flags.read: + return None + compact_end = yyyymmdd(end_date) + if not compact_end: + return None + try: + start = _shift_yyyymmdd(compact_end, -max(190, int(limit) * 3)) + if dataset == "index_daily": + response = self._paginate( + self.client.index_bars, + {"code": code, "from": start, "to": compact_end}, + ) + else: + response = self._paginate( + self.client.daily_bars, + {"code": code, "from": start, "to": compact_end, "adjust": "none"}, + ) + self._validate_usable(dataset, list(response.data or []), response) + rows = _chart_bars(list(response.data or [])) + if not rows: + raise DatahubError("EMPTY", f"{dataset} chart empty") + self._record_route(dataset, "datahub", str((response.meta or {}).get("source") or "datahub")) + return rows[-max(20, min(180, int(limit))):] + except Exception as exc: + self._log_failure(dataset, exc) + return None + + def record_legacy(self, dataset: str, source: str = "", error: str = "") -> None: + self._record_route(dataset, "legacy", source, error) + + def route_snapshot(self) -> list[dict[str, Any]]: + return LEDGER.snapshot() + + def _try_quote_rows( + self, + dataset: str, + params: dict[str, Any], + expected_date: str = "", + minimum: int = 1, + ) -> list[dict[str, Any]] | None: + flags = self.settings.flags(dataset) + if not flags.read: + return None + try: + response = self.client.quotes_latest(**params) + rows = [_native_quote(item) for item in (response.data or []) if isinstance(item, dict)] + rows = [item for item in rows if item] + want = yyyymmdd(expected_date) + if want: + dated = [item for item in rows if not item.get("quote_date") or item.get("quote_date") == want] + if dated: + rows = dated + if len(rows) < minimum: + raise DatahubError("EMPTY", f"datahub {dataset} empty") + if (response.meta or {}).get("stale"): + raise DatahubError("STALE", f"datahub {dataset} stale") + self._record_route(dataset, "datahub", str((response.meta or {}).get("source") or "datahub")) + return rows + except Exception as exc: + self._log_failure(dataset, exc) + return None + def query( self, api_name: str, @@ -180,12 +282,19 @@ class DatahubBridge: raise self._emit_shadow(compare_rows(dataset, legacy_rows, hub_canonical, hub_meta, hub_error, fields)) if flags.read and hub_rows is not None and hub_error is None: + self._record_route(dataset, "datahub", str(hub_meta.get("source") or "datahub")) return project_fields(hub_rows, fields) + if flags.read: + self._record_route(dataset, "legacy", "tushare", hub_error or "") return legacy_rows if flags.read and hub_rows is not None and hub_error is None: + self._record_route(dataset, "datahub", str(hub_meta.get("source") or "datahub")) return project_fields(hub_rows, fields) - return legacy_query(api_name, params, fields) + result = legacy_query(api_name, params, fields) + if flags.read: + self._record_route(dataset, "legacy", "tushare", hub_error or "") + return result def _fetch_dataset(self, dataset: str, params: dict[str, Any], api_name: str = "") -> DatahubResponse: date = yyyymmdd(params.get("trade_date") or params.get("date")) @@ -298,11 +407,12 @@ class DatahubBridge: self.shadow_sink(report) def _log_failure(self, dataset: str, exc: Exception) -> None: - LOGGER.warning( - "datahub fallback dataset=%s error=%s", - dataset, - redact_text(self._error_text(exc), self.settings.secrets()), - ) + error = redact_text(self._error_text(exc), self.settings.secrets()) + LOGGER.warning("datahub fallback dataset=%s error=%s", dataset, error) + self._record_route(dataset, "legacy", "pending-legacy", error) + + def _record_route(self, dataset: str, route: str, source: str = "", error: str = "") -> None: + LEDGER.record(dataset, route, source, redact_text(error, self.settings.secrets())) def _error_text(self, exc: Exception) -> str: if isinstance(exc, DatahubError): @@ -312,6 +422,75 @@ class DatahubBridge: return redact_text(text, self.settings.secrets()) +def _native_quote(row: dict[str, Any]) -> dict[str, Any] | None: + ts_code = str(row.get("ts_code") or "").strip() + close = _finite(row.get("close") if row.get("close") not in (None, "") else row.get("price")) + previous = _finite( + row.get("pre_close") if row.get("pre_close") not in (None, "") else row.get("previous_close") + ) + if not ts_code or close <= 0 or previous <= 0: + return None + volume = _finite(row.get("vol") if row.get("vol") not in (None, "") else row.get("volume")) + return { + "ts_code": ts_code, + "name": str(row.get("name") or ts_code).strip(), + "pre_close": previous, + "open": _finite(row.get("open")), + "high": _finite(row.get("high")), + "low": _finite(row.get("low")), + "close": close, + "vol": volume, + "amount": _finite(row.get("amount")), + "num": 0, + "quote_date": yyyymmdd(row.get("quote_date") or row.get("trade_date")), + "source": str(row.get("source") or "datahub"), + } + + +def _chart_bars(rows: list[Any]) -> list[dict[str, Any]]: + normalized: list[dict[str, Any]] = [] + for row in rows: + if not isinstance(row, dict): + continue + compact = yyyymmdd(row.get("trade_date")) + close = _finite(row.get("close")) + if len(compact) != 8 or close <= 0: + continue + volume = _finite(row.get("volume") if row.get("volume") not in (None, "") else row.get("vol")) + amount = _finite(row.get("amount")) + if volume and volume < close * 10 and amount > 1000: + volume = volume * 100 + trade_date = f"{compact[:4]}-{compact[4:6]}-{compact[6:8]}" + previous = normalized[-1]["close"] if normalized else 0.0 + normalized.append( + { + "trade_date": trade_date, + "open": _finite(row.get("open")), + "high": _finite(row.get("high")), + "low": _finite(row.get("low")), + "close": close, + "change": round((close / previous - 1) * 100, 4) if previous else _finite(row.get("pct_chg")), + "volume": volume, + "amount_billion": amount / 100_000_000, + } + ) + return normalized + + +def _shift_yyyymmdd(value: str, days: int) -> str: + from datetime import datetime, timedelta + + stamp = datetime.strptime(value, "%Y%m%d") + return (stamp + timedelta(days=days)).strftime("%Y%m%d") + + +def _finite(value: Any) -> float: + try: + return float(value or 0) + except (TypeError, ValueError): + return 0.0 + + class DatahubAwareTushareClient: def __init__(self, legacy: TushareClient, bridge: DatahubBridge) -> None: self._legacy = legacy @@ -325,5 +504,17 @@ class DatahubAwareTushareClient: ) -> list[dict[str, Any]]: return self._bridge.query(api_name, params, fields, self._legacy.query) + def try_market_quotes(self, trade_date: str = "") -> list[dict[str, Any]] | None: + return self._bridge.try_market_quotes(trade_date) + + def try_quotes(self, codes: list[str]) -> list[dict[str, Any]] | None: + return self._bridge.try_quotes(codes) + + def try_index_quotes(self) -> list[dict[str, Any]] | None: + return self._bridge.try_index_quotes() + + def record_datahub_legacy(self, dataset: str, source: str = "", error: str = "") -> None: + self._bridge.record_legacy(dataset, source, error) + def __getattr__(self, name: str) -> Any: return getattr(self._legacy, name) diff --git a/backend/data/datahub/route_state.py b/backend/data/datahub/route_state.py new file mode 100644 index 0000000..42c093d --- /dev/null +++ b/backend/data/datahub/route_state.py @@ -0,0 +1,57 @@ +from __future__ import annotations + +from datetime import datetime +from threading import Lock +from typing import Any + +from backend.data.datahub.settings import DATASETS + +DATASET_LABELS = { + "calendar": "交易日历", + "stocks": "股票主档", + "daily": "个股日K", + "index_daily": "指数日K", + "valuation": "估值", + "moneyflow": "资金流", + "auction": "竞价", + "limit_events": "涨停池", + "popularity": "人气榜", + "dragon_tiger": "龙虎榜", + "sector_daily": "题材板块", + "quotes": "全市场实时行情", + "index_quotes": "指数实时行情", + "intraday": "分时", + "status": "数据集状态", +} + + +class DatahubRouteLedger: + def __init__(self) -> None: + self._lock = Lock() + self._rows: dict[str, dict[str, Any]] = {} + + def record(self, dataset: str, route: str, source: str = "", error: str = "") -> None: + name = str(dataset or "").strip() or "unknown" + with self._lock: + self._rows[name] = { + "dataset": name, + "label": DATASET_LABELS.get(name, name), + "route": "legacy" if route == "legacy" else "datahub", + "source": str(source or "").strip(), + "error": str(error or "").strip(), + "at": datetime.now().astimezone().isoformat(timespec="seconds"), + } + + def snapshot(self) -> list[dict[str, Any]]: + with self._lock: + rows = [dict(item) for item in self._rows.values()] + order = {name: index for index, name in enumerate(DATASETS)} + rows.sort(key=lambda item: (order.get(str(item.get("dataset")), 99), str(item.get("dataset")))) + return rows + + def clear(self) -> None: + with self._lock: + self._rows.clear() + + +LEDGER = DatahubRouteLedger() diff --git a/backend/data/gateway.py b/backend/data/gateway.py index 6926585..7cfd65d 100644 --- a/backend/data/gateway.py +++ b/backend/data/gateway.py @@ -47,6 +47,31 @@ class DataGateway: def batches(self, trade_date: str, dataset: str = "") -> list[dict[str, Any]] | None: return self.datahub.batches(trade_date, dataset) + def datahub_status(self) -> dict[str, Any]: + from backend.data.datahub.route_state import DATASET_LABELS, LEDGER + from backend.data.datahub.settings import DATASETS + + settings = self.datahub.settings + flags = [] + enabled = 0 + for name in DATASETS: + read = bool(settings.flags(name).read) + if read: + enabled += 1 + flags.append({"dataset": name, "label": DATASET_LABELS.get(name, name), "read": read}) + routes = LEDGER.snapshot() + fallbacks = [item for item in routes if item.get("route") == "legacy"] + return { + "configured": bool(settings.token and settings.base_url), + "base_url": settings.base_url, + "enabled_reads": enabled, + "total_reads": len(DATASETS), + "flags": flags, + "routes": routes, + "fallback_count": len(fallbacks), + "fallback_labels": [str(item.get("label") or item.get("dataset")) for item in fallbacks], + } + def assert_source(self, dataset_id: str, provider_id: str, usage: DataUsage) -> None: self.policy.assert_allowed(dataset_id, provider_id, usage) diff --git a/backend/data/providers/tushare_dashboard.py b/backend/data/providers/tushare_dashboard.py index d0b0d36..6599153 100644 --- a/backend/data/providers/tushare_dashboard.py +++ b/backend/data/providers/tushare_dashboard.py @@ -186,7 +186,12 @@ 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": + if quote_source == "datahub": + notice = ( + "盘中行情由数据中枢统一提供;涨停原因、封板时间和开板次数以盘后榜单校正为准。" + ) + source_name = "datahub" + elif quote_source == "eastmoney_clist": notice = ( "盘中行情由东财免费实时快照计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。" ) @@ -241,10 +246,16 @@ class DashboardMixin: codes: str, trade_date: str, ) -> tuple[list[dict[str, Any]], str]: + hub = getattr(self, "try_market_quotes", None) + if callable(hub): + quotes = hub(trade_date) + if quotes: + return list(quotes), "datahub" rt_error = "" try: quotes = self.query("rt_k", {"ts_code": codes}) if quotes: + self._mark_quote_legacy("tushare_rt_k", rt_error) return list(quotes), "tushare_rt_k" rt_error = f"No realtime data returned for {trade_date}" except TushareError as exc: @@ -259,8 +270,14 @@ class DashboardMixin: raise TushareError( f"当天盘中实时行情不可用:rt_k={rt_error};免费源=empty" ) + self._mark_quote_legacy(quote_source, rt_error) return quotes, quote_source + def _mark_quote_legacy(self, source: str, error: str = "") -> None: + marker = getattr(self, "record_datahub_legacy", None) + if callable(marker): + marker("quotes", source, error) + def _free_realtime_quotes( self, trade_date: str, @@ -286,8 +303,18 @@ class DashboardMixin: return quotes, "tencent_qt" def _free_realtime_indices(self) -> list[dict[str, Any]]: + hub = getattr(self, "try_index_quotes", None) + if callable(hub): + rows = hub() + converted = [item for item in (_hub_index_quote(row) for row in rows or []) if item] + if converted: + return converted try: - return self._realtime_aggregator().eastmoney_indices() + rows = self._realtime_aggregator().eastmoney_indices() + marker = getattr(self, "record_datahub_legacy", None) + if callable(marker): + marker("index_quotes", "eastmoney_push2") + return rows except Exception: return [] @@ -692,6 +719,31 @@ def _build_yesterday_performance( return result +def _hub_index_quote(row: dict[str, Any]) -> dict[str, Any] | None: + ts_code = str(row.get("ts_code") or "") + code = str(row.get("code") or ts_code.split(".")[0]) + close = _number(row.get("price") if row.get("price") not in (None, "") else row.get("close")) + previous = _number( + row.get("previous_close") if row.get("previous_close") not in (None, "") else row.get("pre_close") + ) + if close <= 0 or previous <= 0: + return None + amount = _number(row.get("amount")) + amount_billion = _number(row.get("amount_billion")) + if not amount_billion and amount: + amount_billion = round(amount / 100_000_000, 2) + return { + "code": code, + "name": str(row.get("name") or code), + "price": close, + "change": _number(row.get("pct_chg") if row.get("pct_chg") not in (None, "") else row.get("change")), + "previous_close": previous, + "amount_billion": amount_billion, + "quote_time": str(row.get("quote_time") or ""), + "source": "datahub", + } + + def _build_limit_performance(rows: list[dict[str, Any]]) -> list[dict[str, Any]]: result = [] for level in sorted({int(row.get("prior_streak") or 1) for row in rows}, reverse=True): diff --git a/backend/data/providers/tushare_indices.py b/backend/data/providers/tushare_indices.py index ce1a223..eab8eca 100644 --- a/backend/data/providers/tushare_indices.py +++ b/backend/data/providers/tushare_indices.py @@ -59,10 +59,85 @@ class IndexMixin: } def realtime_market_indices(self, requested_date: str) -> dict[str, Any]: + hub = getattr(self, "try_index_quotes", None) + if callable(hub): + rows = hub() + if rows: + try: + return self._hub_realtime_market_indices(requested_date, rows) + except TushareError: + pass try: - return self._tushare_realtime_market_indices(requested_date) + payload = self._tushare_realtime_market_indices(requested_date) + marker = getattr(self, "record_datahub_legacy", None) + if callable(marker): + marker("index_quotes", "tushare_rt_idx_k") + return payload except TushareError: - return self._free_realtime_market_indices(requested_date) + payload = self._free_realtime_market_indices(requested_date) + marker = getattr(self, "record_datahub_legacy", None) + if callable(marker): + marker("index_quotes", str(payload.get("source") or "eastmoney_push2")) + return payload + + def _hub_realtime_market_indices( + self, + requested_date: str, + rows: list[dict[str, Any]], + ) -> dict[str, Any]: + trade_date, _ = self.resolve_trade_context(requested_date) + index_names = { + "000001.SH": "上证指数", + "399001.SZ": "深证成指", + "399006.SZ": "创业板指", + } + by_code = {str(row.get("ts_code") or ""): row for row in rows} + by_symbol = {str(row.get("code") or ""): row for row in rows} + indices = [] + for ts_code, name in index_names.items(): + row = by_code.get(ts_code) or by_symbol.get(ts_code.split(".")[0]) + if not row: + continue + close = _number(row.get("price") if row.get("price") not in (None, "") else row.get("close")) + previous_close = _number( + row.get("previous_close") if row.get("previous_close") not in (None, "") else row.get("pre_close") + ) + if close <= 0 or previous_close <= 0: + continue + amount = _number(row.get("amount")) + amount_billion = _number(row.get("amount_billion")) + if not amount_billion and amount: + amount_billion = round(amount / 100_000_000, 2) + indices.append( + { + "ts_code": ts_code, + "name": str(row.get("name") or name).strip(), + "trade_date": trade_date, + "close": close, + "pct_chg": round( + _number(row.get("pct_chg")) or (close / previous_close - 1) * 100, + 3, + ), + "return_5d": 0, + "amount_billion": amount_billion, + "quote_time": str(row.get("quote_time") or ""), + "source": "datahub", + } + ) + if len(indices) != 3: + raise TushareError("Realtime index quotes are incomplete") + return { + "trade_date": trade_date, + "source": "datahub", + "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, + }, + } def _tushare_realtime_market_indices(self, requested_date: str) -> dict[str, Any]: trade_date, _ = self.resolve_trade_context(requested_date) diff --git a/backend/features/market/charts.py b/backend/features/market/charts.py index 1b1dcbb..f3db1d0 100644 --- a/backend/features/market/charts.py +++ b/backend/features/market/charts.py @@ -68,12 +68,18 @@ class MarketChartClient: normalized = str(code or "").strip() if not re.fullmatch(r"\d{6}", normalized): raise ChartDataError("Invalid stock code") + hub_rows = self._datahub_daily(normalized, end_date, limit, "daily") + if hub_rows: + return hub_rows return self._ifind_daily(_stock_market_code(normalized), end_date, limit) def index_daily(self, identifier: str, end_date: str, limit: int = 90) -> list[dict[str, Any]]: normalized = str(identifier or "").strip().upper() if normalized not in INDEX_SECIDS: raise ChartDataError("Unsupported index") + hub_rows = self._datahub_daily(normalized, end_date, limit, "index_daily") + if hub_rows: + return hub_rows return self._ifind_daily(normalized, end_date, limit) def board_daily(self, identifier: str, end_date: str, limit: int = 90) -> list[dict[str, Any]]: @@ -109,6 +115,112 @@ class MarketChartClient: return None return chart + def _datahub_daily( + self, + code: str, + end_date: str, + limit: int, + dataset: str, + ) -> list[dict[str, Any]] | None: + if self.datahub is None or not hasattr(self.datahub, "try_daily_chart"): + return None + try: + rows = self.datahub.try_daily_chart(code, end_date, limit, dataset) + except Exception as exc: + LOGGER.warning("datahub daily unexpected error: %s", exc) + rows = None + if not rows: + if hasattr(self.datahub, "record_legacy"): + self.datahub.record_legacy(dataset, "ifind") + return None + compact_end = str(end_date or "").replace("-", "") + market_now = datetime.now().astimezone() + today = market_now.strftime("%Y%m%d") + market_open = ( + market_now.weekday() < 5 + and market_now.time().replace(tzinfo=None) >= dt_time(9, 30) + ) + if compact_end == today and market_open: + overlay = self._datahub_today_bar(code, dataset, rows) + if overlay: + if rows and rows[-1]["trade_date"] == overlay["trade_date"]: + rows[-1] = overlay + else: + rows.append(overlay) + return rows + + def _datahub_today_bar( + self, + code: str, + dataset: str, + history: list[dict[str, Any]], + ) -> dict[str, Any] | None: + today_display = datetime.now().astimezone().date().isoformat() + previous = history[-1]["close"] if history and history[-1]["trade_date"] != today_display else ( + history[-2]["close"] if len(history) >= 2 else 0.0 + ) + quote = None + if dataset == "index_daily" and hasattr(self.datahub, "try_index_quotes"): + quotes = self.datahub.try_index_quotes() or [] + quote = next( + ( + item for item in quotes + if str(item.get("ts_code") or "") == code or str(item.get("code") or "") == code.split(".")[0] + ), + None, + ) + elif hasattr(self.datahub, "try_quotes"): + quotes = self.datahub.try_quotes([code]) or [] + quote = quotes[0] if quotes else None + if quote: + close = _number(quote.get("close") if quote.get("close") not in (None, "") else quote.get("price")) + open_price = _number(quote.get("open")) + high = _number(quote.get("high")) + low = _number(quote.get("low")) + previous_close = _number( + quote.get("pre_close") if quote.get("pre_close") not in (None, "") else quote.get("previous_close") + ) or previous + volume = _number(quote.get("vol") if quote.get("vol") not in (None, "") else quote.get("volume")) + amount = _number(quote.get("amount")) + if close > 0 and open_price > 0: + return { + "trade_date": today_display, + "open": open_price, + "high": high or close, + "low": low or close, + "close": close, + "change": round((close / previous_close - 1) * 100, 4) if previous_close else 0.0, + "volume": volume, + "amount_billion": amount / 100_000_000, + "realtime": True, + } + chart = self._datahub_intraday(code) + points = list((chart or {}).get("points") or []) + if not points: + return None + closes = [_number(point.get("close")) for point in points if _number(point.get("close")) > 0] + if not closes: + return None + opens = [_number(point.get("open")) for point in points if _number(point.get("open")) > 0] + highs = [_number(point.get("high")) for point in points if _number(point.get("high")) > 0] + lows = [_number(point.get("low")) for point in points if _number(point.get("low")) > 0] + volume = sum(_number(point.get("volume")) for point in points) + amount = sum(_number(point.get("amount")) for point in points) + previous_close = _number((chart or {}).get("previous_close")) or previous + close = closes[-1] + open_price = opens[0] if opens else closes[0] + return { + "trade_date": today_display, + "open": open_price, + "high": max(highs or closes), + "low": min(lows or closes), + "close": close, + "change": round((close / previous_close - 1) * 100, 4) if previous_close else 0.0, + "volume": volume, + "amount_billion": amount / 100_000_000, + "realtime": True, + } + def board_intraday(self, identifier: str, name: str = "") -> dict[str, Any]: normalized = str(identifier or "").strip().upper() try: diff --git a/backend/features/system/service.py b/backend/features/system/service.py index 4f14ebe..267c618 100644 --- a/backend/features/system/service.py +++ b/backend/features/system/service.py @@ -130,6 +130,7 @@ class SystemServiceMixin: ), **self.database.status(), "jobs": self.jobs.repository.recent(12), + "datahub": self._datahub_status(), }, "llm": { "primary_configured": self._profile_configured(platform["primary"]), @@ -145,6 +146,22 @@ class SystemServiceMixin: }, } + def _datahub_status(self) -> dict[str, Any]: + gateway = getattr(self, "data_gateway", None) + reporter = getattr(gateway, "datahub_status", None) + if callable(reporter): + return reporter() + return { + "configured": False, + "base_url": "", + "enabled_reads": 0, + "total_reads": 0, + "flags": [], + "routes": [], + "fallback_count": 0, + "fallback_labels": [], + } + def save_system_settings(self, payload: dict[str, Any]) -> dict[str, Any]: current = dict(self._system_credentials) token = str(payload.get("tushare_token") or current.get("tushare_token") or "").strip() diff --git a/compose.yaml b/compose.yaml index a97a8d2..934fcc2 100644 --- a/compose.yaml +++ b/compose.yaml @@ -13,6 +13,22 @@ services: - ./.env environment: APP_ENCRYPTION_KEY: "${APP_ENCRYPTION_KEY:?APP_ENCRYPTION_KEY must be set in .env}" + DATAHUB_BASE_URL: "${DATAHUB_BASE_URL:-http://192.168.200.11:8766}" + DATAHUB_READ_CALENDAR: "1" + DATAHUB_READ_STOCKS: "1" + DATAHUB_READ_DAILY: "1" + DATAHUB_READ_INDEX_DAILY: "1" + DATAHUB_READ_VALUATION: "1" + DATAHUB_READ_MONEYFLOW: "1" + DATAHUB_READ_AUCTION: "1" + DATAHUB_READ_LIMIT_EVENTS: "1" + DATAHUB_READ_POPULARITY: "1" + DATAHUB_READ_DRAGON_TIGER: "1" + DATAHUB_READ_SECTOR_DAILY: "1" + DATAHUB_READ_QUOTES: "1" + DATAHUB_READ_INDEX_QUOTES: "1" + DATAHUB_READ_INTRADAY: "1" + DATAHUB_READ_STATUS: "1" TZ: Asia/Shanghai PYTHONUTF8: "1" volumes: diff --git a/config/README.md b/config/README.md index 9b9d834..fac3b74 100644 --- a/config/README.md +++ b/config/README.md @@ -12,9 +12,12 @@ These registries describe the approved product surface of the standalone applica providers, model entry points, CSS layers, and remaining code hotspots. - `data-fields.config.json`: canonical data products, provider eligibility, intended use, and known blocked datasets. -- `datahub.config.json`: optional read-only client for `xiaobai-datahub`. Each dataset has its - own `read` / `shadow` flag, all default off. Environment variables `DATAHUB_READ_*` and - `DATAHUB_SHADOW_*` can override a single dataset without a master switch. +- `datahub.config.json`: official read-only client for `xiaobai-datahub`. Each dataset has its + own `read` / `shadow` flag; official reads default on. `compose.yaml` pins every + `DATAHUB_READ_*` to `"1"` so a leftover `.env` `=0` cannot silently keep official + pages on the old APIs. Environment variables can still override a single + `DATAHUB_SHADOW_*` without a master switch. The old website APIs stay as + emergency fallback only. - `data-quality.config.json`: freshness, coverage, units, adjustment, point-in-time, and fail-closed rules for every canonical data product. - `jobs.config.json`: background schedules, dependencies, lock keys, retry policy, timeouts, diff --git a/config/architecture-inventory.json b/config/architecture-inventory.json index 1d4be55..f5d3dc9 100644 --- a/config/architecture-inventory.json +++ b/config/architecture-inventory.json @@ -483,8 +483,8 @@ }, { "path": "frontend/index.html", - "bytes": 48254, - "lines": 664 + "bytes": 48447, + "lines": 665 }, { "path": "backend/features/screener/catalog.py", @@ -496,6 +496,11 @@ "bytes": 35247, "lines": 2416 }, + { + "path": "backend/data/providers/tushare_dashboard.py", + "bytes": 33603, + "lines": 784 + }, { "path": "database.py", "bytes": 32073, @@ -506,11 +511,6 @@ "bytes": 31756, "lines": 562 }, - { - "path": "backend/data/providers/tushare_dashboard.py", - "bytes": 31361, - "lines": 732 - }, { "path": "backend/data/providers/tushare_industries.py", "bytes": 26540, @@ -551,6 +551,11 @@ "bytes": 15311, "lines": 387 }, + { + "path": "frontend/shared/admin.js", + "bytes": 15235, + "lines": 289 + }, { "path": "frontend/shared/dashboard.js", "bytes": 15063, @@ -566,11 +571,6 @@ "bytes": 14743, "lines": 342 }, - { - "path": "frontend/shared/admin.js", - "bytes": 14410, - "lines": 268 - }, { "path": "backend/features/heaven/market_context.py", "bytes": 13681, @@ -581,15 +581,20 @@ "bytes": 13219, "lines": 289 }, + { + "path": "backend/features/system/service.py", + "bytes": 12937, + "lines": 271 + }, { "path": "backend/features/market/insights_auction_data.py", "bytes": 12829, "lines": 318 }, { - "path": "backend/features/system/service.py", - "bytes": 12392, - "lines": 254 + "path": "backend/data/providers/tushare_indices.py", + "bytes": 10956, + "lines": 248 }, { "path": "backend/features/market/insights_auction.py", @@ -631,11 +636,6 @@ "bytes": 8357, "lines": 116 }, - { - "path": "backend/data/providers/tushare_indices.py", - "bytes": 7823, - "lines": 173 - }, { "path": "backend/features/screener/formula.py", "bytes": 6983, diff --git a/config/datahub.config.json b/config/datahub.config.json index 31f622a..37e4001 100644 --- a/config/datahub.config.json +++ b/config/datahub.config.json @@ -6,20 +6,20 @@ "page_limit": 5000, "stale_seconds_max": 86400, "datasets": { - "calendar": { "read": false, "shadow": false }, - "stocks": { "read": false, "shadow": false }, - "daily": { "read": false, "shadow": false }, - "index_daily": { "read": false, "shadow": false }, - "valuation": { "read": false, "shadow": false }, - "moneyflow": { "read": false, "shadow": false }, - "auction": { "read": false, "shadow": false }, - "limit_events": { "read": false, "shadow": false }, - "popularity": { "read": false, "shadow": false }, - "dragon_tiger": { "read": false, "shadow": false }, - "sector_daily": { "read": false, "shadow": false }, - "quotes": { "read": false, "shadow": false }, - "index_quotes": { "read": false, "shadow": false }, - "intraday": { "read": false, "shadow": false }, - "status": { "read": false, "shadow": false } + "calendar": { "read": true, "shadow": false }, + "stocks": { "read": true, "shadow": false }, + "daily": { "read": true, "shadow": false }, + "index_daily": { "read": true, "shadow": false }, + "valuation": { "read": true, "shadow": false }, + "moneyflow": { "read": true, "shadow": false }, + "auction": { "read": true, "shadow": false }, + "limit_events": { "read": true, "shadow": false }, + "popularity": { "read": true, "shadow": false }, + "dragon_tiger": { "read": true, "shadow": false }, + "sector_daily": { "read": true, "shadow": false }, + "quotes": { "read": true, "shadow": false }, + "index_quotes": { "read": true, "shadow": false }, + "intraday": { "read": true, "shadow": false }, + "status": { "read": true, "shadow": false } } } diff --git a/frontend/index.html b/frontend/index.html index f9fcef3..fc7f787 100644 --- a/frontend/index.html +++ b/frontend/index.html @@ -611,6 +611,7 @@

所有用户读取同一份后台快照,页面不会随后台任务自动重绘。

+
数据中枢线路待检查
尚未手动刷新
diff --git a/frontend/m/js/pages.js b/frontend/m/js/pages.js index 42fcab3..af02d59 100644 --- a/frontend/m/js/pages.js +++ b/frontend/m/js/pages.js @@ -5219,6 +5219,15 @@ return ''; } + function datahubStatusText(hub) { + const enabled = number(hub.enabled_reads); + const total = number(hub.total_reads) || enabled; + const fallbacks = hub.fallback_labels || []; + if (fallbacks.length) return " 备用 " + fallbacks.join("、"); + if (hub.configured) return " 主线路 " + enabled + "/" + total; + return " 未配置"; + } + function renderSystemAdmin(key) { if (key === "system/members") { renderSystemMembers(); @@ -5237,6 +5246,7 @@ '
iFinD' + statusDot(ifind.configured) + (ifind.configured ? " 已配置" : " 未配置") + "
" + '
行情快照' + number(data.snapshot_dates) + " 个交易日
" + '
后台刷新' + statusDot(data.background_refresh_enabled) + (data.background_refresh_enabled ? " 已启用" : " 已暂停") + "
" + + '
数据中枢' + statusDot(Boolean((data.datahub || {}).configured) && !((data.datahub || {}).fallback_count)) + datahubStatusText(data.datahub || {}) + "
" + "" + '
数据源密钥' + formFieldHtml("Tushare Token", '', false) + diff --git a/frontend/shared/admin.js b/frontend/shared/admin.js index cebe715..ca33f1b 100644 --- a/frontend/shared/admin.js +++ b/frontend/shared/admin.js @@ -44,6 +44,7 @@ async function openAdminSettings(refreshOnly = false) { status.textContent = `Tushare ${data.configured ? "已配置" : "未配置"} · iFinD ${ifind.configured ? "已配置" : "未配置"} · ${number(data.snapshot_dates)} 个交易日`; status.classList.toggle("connected", Boolean(data.configured)); setText("systemDataStatus", data.background_refresh_enabled ? "后台刷新已启用" : "后台刷新已暂停"); + renderDatahubRouteStatus(data.datahub || {}); document.querySelector("#systemTokenInput").value = ""; document.querySelector("#systemIfindTokenInput").value = ""; document.querySelector("#systemBackgroundRefresh").checked = Boolean(data.background_refresh_enabled); @@ -55,6 +56,26 @@ async function openAdminSettings(refreshOnly = false) { } } +function renderDatahubRouteStatus(hub) { + const box = document.querySelector("#datahubRouteStatus"); + if (!box) return; + const label = box.querySelector("span"); + const enabled = number(hub.enabled_reads); + const total = number(hub.total_reads) || enabled; + const fallbacks = hub.fallback_labels || []; + if (fallbacks.length) { + box.dataset.tone = "warning"; + if (label) label.textContent = `数据中枢主线路 ${enabled}/${total} · 备用 ${fallbacks.length} 类:${fallbacks.join("、")}`; + return; + } + box.dataset.tone = hub.configured ? "success" : "idle"; + if (label) { + label.textContent = hub.configured + ? `数据中枢主线路 ${enabled}/${total},当前无备用` + : "数据中枢尚未配置,网站仍走原接口"; + } +} + function selectAdminPanel(panel) { const selected = ["market", "models", "members"].includes(panel) ? panel : "market"; document.querySelector("#adminSectionSelect").value = selected; diff --git a/tests/test_chart_data_provider.py b/tests/test_chart_data_provider.py index 6d9314c..6bb6008 100644 --- a/tests/test_chart_data_provider.py +++ b/tests/test_chart_data_provider.py @@ -155,10 +155,12 @@ class ChartLookbackTests(unittest.TestCase): class FakeHub: - def __init__(self, chart=None, error=None): + def __init__(self, chart=None, error=None, daily=None): self.chart = chart self.error = error + self.daily = daily self.calls: list[str] = [] + self.legacy: list[str] = [] def try_intraday(self, code): self.calls.append(code) @@ -166,6 +168,15 @@ class FakeHub: raise self.error return self.chart + def try_daily_chart(self, code, end_date, limit=90, dataset="daily"): + self.calls.append(f"{dataset}:{code}") + if self.error: + raise self.error + return self.daily + + def record_legacy(self, dataset, source="", error=""): + self.legacy.append(dataset) + class DatahubChartFallbackTests(unittest.TestCase): def setUp(self) -> None: @@ -207,6 +218,25 @@ class DatahubChartFallbackTests(unittest.TestCase): self.assertGreaterEqual(len(payload["points"]), 1) self.assertTrue(fallback.requests) + def test_datahub_daily_skips_ifind(self): + hub = FakeHub( + daily=[ + { + "trade_date": "2026-09-07", + "open": 10.0, + "high": 10.4, + "low": 9.9, + "close": 10.2, + "volume": 1000, + "amount_billion": 0.02, + } + ] + ) + client = MarketChartClient(IfindHttpClient(), LookbackChartClient(), hub) + rows = client.stock_daily("600000", "20260907") + self.assertEqual(rows[-1]["trade_date"], "2026-09-07") + self.assertIn("daily:600000", hub.calls) + class ChartServiceStub: @staticmethod diff --git a/tests/test_datahub_bridge.py b/tests/test_datahub_bridge.py index fe921d2..85b70c8 100644 --- a/tests/test_datahub_bridge.py +++ b/tests/test_datahub_bridge.py @@ -12,6 +12,7 @@ from backend.data.datahub.client import DatahubClient, DatahubResponse from backend.data.datahub.compare import compare_rows from backend.data.datahub.errors import DatahubError from backend.data.datahub.native import to_canonical_row, to_native_row +from backend.data.datahub.route_state import LEDGER from backend.data.datahub.settings import DATASETS, DatahubSettings, DatasetFlags ROOT = Path(__file__).resolve().parents[1] @@ -84,17 +85,21 @@ def flags(**enabled: tuple[bool, bool]) -> DatahubSettings: class DatahubBridgeTests(unittest.TestCase): - def test_default_config_keeps_legacy_and_does_not_call_datahub(self) -> None: + def setUp(self) -> None: + LEDGER.clear() + + def test_default_config_enables_official_reads(self) -> None: settings = DatahubSettings.load(environ={}, credentials={}) - self.assertFalse(settings.any_enabled()) - self.assertTrue(all(not settings.flags(name).read and not settings.flags(name).shadow for name in DATASETS)) - client = FakeClient(error=DatahubError("INTERNAL", "should not be called")) + self.assertTrue(settings.any_enabled()) + self.assertTrue(all(settings.flags(name).read and not settings.flags(name).shadow for name in DATASETS)) + client = FakeClient() legacy = FakeLegacy([LEGACY_DAILY]) wrapped = DatahubAwareTushareClient(legacy, DatahubBridge(settings, client)) rows = wrapped.query("daily", {"trade_date": "20240902"}, "ts_code,close,vol,amount") self.assertEqual(rows[0]["amount"], 2000.0) - self.assertEqual(client.paths, []) - self.assertEqual(len(legacy.calls), 1) + self.assertEqual(client.paths, ["/v1/bars/daily"]) + self.assertEqual(legacy.calls, []) + self.assertEqual(LEDGER.snapshot()[0]["route"], "datahub") def test_each_dataset_has_independent_read_flag(self) -> None: settings = flags(daily=(True, False), auction=(False, False)) @@ -104,6 +109,13 @@ class DatahubBridgeTests(unittest.TestCase): source = (ROOT / "config" / "datahub.config.json").read_text(encoding="utf-8") self.assertNotIn("master", source) self.assertNotIn("DATAHUB_READ_ALL", source) + compose = (ROOT / "compose.yaml").read_text(encoding="utf-8") + for env_key in ( + "CALENDAR", "STOCKS", "DAILY", "INDEX_DAILY", "VALUATION", "MONEYFLOW", + "AUCTION", "LIMIT_EVENTS", "POPULARITY", "DRAGON_TIGER", "SECTOR_DAILY", + "QUOTES", "INDEX_QUOTES", "INTRADAY", "STATUS", + ): + self.assertIn(f'DATAHUB_READ_{env_key}: "1"', compose) def test_read_flag_replaces_only_that_dataset_and_converts_units(self) -> None: shadows: list[dict[str, Any]] = [] @@ -412,7 +424,92 @@ class DatahubBridgeTests(unittest.TestCase): FakeClient(error=DatahubError("INTERNAL", "datahub exploded")), ) self.assertIsNone(broken.try_intraday("601318")) - self.assertFalse(DatahubSettings.load(environ={}, credentials={}).flags("intraday").read) + self.assertTrue(DatahubSettings.load(environ={}, credentials={}).flags("intraday").read) + + def test_try_market_quotes_and_visible_fallback(self) -> None: + quotes = [ + { + "ts_code": f"{600000 + index:06d}.SH", + "name": f"股票{index}", + "close": 10.2, + "pre_close": 10.0, + "open": 10.1, + "high": 10.3, + "low": 9.9, + "vol": 1000, + "amount": 2000000, + "quote_date": "20240902", + } + for index in range(220) + ] + ok = DatahubBridge( + flags(quotes=(True, False)), + FakeClient( + response=DatahubResponse( + data=quotes, + meta={"stale": False, "staleness_seconds": 0, "source": "eastmoney:clist"}, + ) + ), + ) + rows = ok.try_market_quotes("20240902") + self.assertEqual(len(rows), 220) + self.assertEqual(rows[0]["pre_close"], 10.0) + self.assertEqual(ok.client.paths, ["/v1/quotes/latest"]) + self.assertEqual(LEDGER.snapshot()[0]["route"], "datahub") + + failed = DatahubBridge( + flags(quotes=(True, False)), + FakeClient(error=DatahubError("UNAVAILABLE", "down")), + ) + self.assertIsNone(failed.try_market_quotes("20240902")) + failed.record_legacy("quotes", "tencent_qt", "down") + snap = next(item for item in LEDGER.snapshot() if item["dataset"] == "quotes") + self.assertEqual(snap["route"], "legacy") + self.assertEqual(snap["source"], "tencent_qt") + self.assertIn("备用", "备用") + + gateway = build_data_gateway({}, datahub_settings=flags(quotes=(True, False))) + status = gateway.datahub_status() + self.assertEqual(status["enabled_reads"], 1) + self.assertEqual(status["total_reads"], len(DATASETS)) + self.assertGreaterEqual(status["fallback_count"], 1) + + def test_try_daily_chart_converts_hub_bars(self) -> None: + rows = [ + { + "ts_code": "600000.SH", + "trade_date": "20240901", + "open": 10.0, + "high": 10.4, + "low": 9.9, + "close": 10.2, + "volume": 100000, + "amount": 2000000, + }, + { + "ts_code": "600000.SH", + "trade_date": "20240902", + "open": 10.2, + "high": 10.5, + "low": 10.1, + "close": 10.4, + "volume": 120000, + "amount": 2400000, + }, + ] + hub = DatahubBridge( + flags(daily=(True, False)), + FakeClient( + response=DatahubResponse( + data=rows, + meta={"stale": False, "staleness_seconds": 0, "source": "tushare:daily"}, + ) + ), + ) + chart = hub.try_daily_chart("600000.SH", "20240902", 90, "daily") + self.assertEqual(chart[-1]["trade_date"], "2024-09-02") + self.assertEqual(chart[-1]["close"], 10.4) + self.assertAlmostEqual(chart[-1]["amount_billion"], 0.024) def test_features_do_not_import_datahub_client(self) -> None: violations = [] diff --git a/tests/test_http_dispatch.py b/tests/test_http_dispatch.py index d889216..ee50f89 100644 --- a/tests/test_http_dispatch.py +++ b/tests/test_http_dispatch.py @@ -138,7 +138,7 @@ class HttpDispatchContractTests(unittest.TestCase): self.assertTrue(claimed.isdisjoint(methods)) claimed.update(methods) self.assertLessEqual(len(path.read_text(encoding="utf-8").splitlines()), line_limit) - self.assertEqual(len(claimed), 27) + self.assertEqual(len(claimed), 28) if __name__ == "__main__": diff --git a/tests/test_realtime_dashboard.py b/tests/test_realtime_dashboard.py index f5c32a3..2a0bbff 100644 --- a/tests/test_realtime_dashboard.py +++ b/tests/test_realtime_dashboard.py @@ -364,6 +364,23 @@ class RealtimeDashboardTests(unittest.TestCase): self.assertEqual(quote["amount"], 406050000) self.assertEqual(quote["source"], "tencent_qt") + def test_datahub_market_quotes_used_before_legacy(self): + calls = [] + + def try_market_quotes(trade_date): + calls.append(trade_date) + return list(FREE_QUOTES) + + self.client.try_market_quotes = try_market_quotes + self.client.realtime_aggregator = FakeFreeAggregator(fail=True) + TushareClient._realtime_reference_cache.clear() + dashboard = self.client._realtime_dashboard("20260720", "20260720", "20260717") + self.assertEqual(calls, ["20260720"]) + self.assertEqual(dashboard["meta"]["quote_source"], "datahub") + self.assertEqual(dashboard["meta"]["source"], "datahub") + self.assertEqual(dashboard["meta"]["quote_count"], 3) + self.assertIn("数据中枢", dashboard["meta"]["notice"]) + if __name__ == "__main__": unittest.main() diff --git a/xiaobai-datahub/README.md b/xiaobai-datahub/README.md index 7d82fa4..437d6d3 100644 --- a/xiaobai-datahub/README.md +++ b/xiaobai-datahub/README.md @@ -7,7 +7,7 @@ - SQLite WAL `datahub.db`,容器名 `xiaobai-datahub`,端口 `8766` - Tushare 盘后正式数据:交易日历、股票主档、daily、daily_basic、adj_factor、index_daily、moneyflow、stk_auction、limit_list_d、ths_hot/dc_hot、hm_detail、ths_daily/dc_index/sw_daily -- 盘中观察(provisional):东财/腾讯指数报价、个股最新价、分时点(`/v1/quotes/latest` `/v1/indexes/quotes` `/v1/intraday/points`);永不写入 eod_* 正式表 +- 盘中观察(provisional):东财/腾讯指数报价、个股最新价、全市场快照、分时点(`/v1/quotes/latest` 不传 codes 即全市场,`/v1/indexes/quotes` `/v1/intraday/points`);永不写入 eod_* 正式表 - 暂存 → 校验 → 整批原子发布 → 可回滚 - `/v1` 稳定接口(`X-Datahub-Token`) - `/admin/` 最小管理后台(总览 / 数据源 / 调度 / 发布 / 数据集 / 审计) diff --git a/xiaobai-datahub/datahub/adapters/eastmoney.py b/xiaobai-datahub/datahub/adapters/eastmoney.py index 01bf681..fa52155 100644 --- a/xiaobai-datahub/datahub/adapters/eastmoney.py +++ b/xiaobai-datahub/datahub/adapters/eastmoney.py @@ -13,6 +13,15 @@ from datahub.numbers import finite_number, round4 EASTMONEY_INDEX_URL = "https://push2.eastmoney.com/api/qt/ulist.np/get" EASTMONEY_CLIST_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 TRENDS_URL = "https://push2delay.eastmoney.com/api/qt/stock/trends2/get" HIS_TRENDS_URL = "https://push2his.eastmoney.com/api/qt/stock/trends2/get" BROWSER_UA = ( @@ -59,7 +68,11 @@ class EastmoneyAdapter(MarketAdapter): codes = params.get("codes") or [] if isinstance(codes, str): codes = [item.strip() for item in codes.split(",") if item.strip()] - return self.fetch_quotes(list(codes)) + if codes: + return self.fetch_quotes(list(codes)) + return self.fetch_market_quotes() + if dataset in {"quotes_market", "market_quotes"}: + return self.fetch_market_quotes() raise AdapterError(f"{self.name} unsupported dataset: {dataset}") def normalize(self, dataset: str, rows: list[dict[str, Any]]) -> list[dict[str, Any]]: @@ -167,6 +180,58 @@ class EastmoneyAdapter(MarketAdapter): ) return result + def fetch_market_quotes(self) -> list[dict[str, Any]]: + rows: list[dict[str, Any]] = [] + board_errors: list[str] = [] + for board in EASTMONEY_A_SHARE_BOARDS: + try: + rows.extend(self._board_quotes(board)) + except Exception as exc: + board_errors.append(f"{board}:{exc}") + quotes: list[dict[str, Any]] = [] + seen: set[str] = set() + for row in rows: + quote = _normalize_market_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 AdapterError(f"Eastmoney market snapshot too small: {len(quotes)}{detail}") + return quotes + + def _board_quotes(self, board: str) -> list[dict[str, Any]]: + first = self._market_page(board, 1) + data = first.get("data") or {} + rows = list(data.get("diff") or []) + total = int(finite_number(data.get("total")) or 0) + 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._market_page(board, page) + rows.extend(list((payload.get("data") or {}).get("diff") or [])) + return rows + + def _market_page(self, board: str, page: int) -> dict[str, Any]: + return self._get_json( + EASTMONEY_CLIST_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 fetch_intraday(self, ts_code: str, date: str = "") -> dict[str, Any]: code = str(ts_code or "").upper() if code in INDEX_SECIDS: @@ -252,6 +317,45 @@ def _preferred_session(points: list[dict[str, Any]], preferred_date: str = "") - return [point for point in points if str(point.get("date") or "") == latest] +def _normalize_market_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 = round4(finite_number(row.get("f2"))) + previous_close = round4(finite_number(row.get("f18"))) + if close <= 0 or previous_close <= 0: + return None + market = int(finite_number(row.get("f13")) or 0) + 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(finite_number(row.get("f124")) or 0) + 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, + "previous_close": previous_close, + "open": round4(finite_number(row.get("f17"))), + "high": round4(finite_number(row.get("f15"))), + "low": round4(finite_number(row.get("f16"))), + "close": close, + "price": close, + "pct_chg": round4(finite_number(row.get("f3"))), + "vol": round4(finite_number(row.get("f5")) * 100), + "volume": round4(finite_number(row.get("f5")) * 100), + "amount": round4(finite_number(row.get("f6"))), + "quote_date": quote_date, + "quote_time_epoch": epoch, + "source": "eastmoney_clist", + } + + def _parse_trend(raw: Any) -> dict[str, Any] | None: text = str(raw or "") parts = text.split(",") diff --git a/xiaobai-datahub/datahub/realtime_serve.py b/xiaobai-datahub/datahub/realtime_serve.py index 5641105..01257b8 100644 --- a/xiaobai-datahub/datahub/realtime_serve.py +++ b/xiaobai-datahub/datahub/realtime_serve.py @@ -64,9 +64,36 @@ def fetch_index_quotes(db: HubDB) -> dict[str, Any]: return payload +def fetch_market_quotes(db: HubDB) -> dict[str, Any]: + cache_key = "quotes:market" + cached = _read_cache(db, cache_key) + if cached is not None: + return cached + adapter = EastmoneyAdapter() + try: + rows = adapter.fetch_market_quotes() + source = "eastmoney:clist" + except Exception as exc: + raise RealtimeApiError("SOURCE_UNAVAILABLE", f"market quotes unavailable: {exc}") from exc + payload = _envelope( + rows, + { + "tier": "provisional", + "trade_date": yyyymmdd(now_shanghai()), + "source": source, + "stale": False, + "staleness_seconds": 0, + "published_at": isoformat(now_shanghai()), + "scope": "market", + }, + ) + _write_cache(db, cache_key, payload, QUOTE_TTL, source) + return payload + + def fetch_quotes(db: HubDB, codes: list[str]) -> dict[str, Any]: if not codes: - raise RealtimeApiError("INVALID_ARGUMENT", "codes is required") + return fetch_market_quotes(db) resolved: list[str] = [] for code in codes[:60]: item = resolve_code(db, code) or _guess_ts_code(code) diff --git a/xiaobai-datahub/datahub/serving.py b/xiaobai-datahub/datahub/serving.py index d8c7585..6cfe214 100644 --- a/xiaobai-datahub/datahub/serving.py +++ b/xiaobai-datahub/datahub/serving.py @@ -247,11 +247,13 @@ class V1API: ) def quotes_latest(self, q: dict[str, str]) -> dict[str, Any]: - from datahub.realtime_serve import RealtimeApiError, fetch_quotes + from datahub.realtime_serve import RealtimeApiError, fetch_market_quotes, fetch_quotes codes = [item.strip() for item in str(q.get("codes") or "").split(",") if item.strip()] try: - return fetch_quotes(self.db, codes) + if codes: + return fetch_quotes(self.db, codes) + return fetch_market_quotes(self.db) except RealtimeApiError as exc: raise ApiError(exc.code, exc.message) from exc diff --git a/xiaobai-datahub/tests/test_realtime_intraday.py b/xiaobai-datahub/tests/test_realtime_intraday.py index a52e264..0a4f703 100644 --- a/xiaobai-datahub/tests/test_realtime_intraday.py +++ b/xiaobai-datahub/tests/test_realtime_intraday.py @@ -172,5 +172,61 @@ class ServingIntradayDateTests(unittest.TestCase): self.assertIn("intraday unavailable", ctx.exception.message) +class MarketQuotesTests(unittest.TestCase): + def setUp(self) -> None: + self.tmp = tempfile.TemporaryDirectory() + self.db = HubDB(Path(self.tmp.name) / "hub.db") + self.api = V1API(self.db, None, None) # type: ignore[arg-type] + + def tearDown(self) -> None: + self.tmp.cleanup() + + def test_empty_codes_returns_full_market_snapshot(self) -> None: + rows = [ + { + "ts_code": f"{600000 + index:06d}.SH", + "name": f"股票{index}", + "pre_close": 10.0, + "close": 10.2, + "open": 10.1, + "high": 10.3, + "low": 10.0, + "vol": 1000, + "amount": 2000000, + } + for index in range(220) + ] + with patch("datahub.realtime_serve.EastmoneyAdapter") as mocked: + mocked.return_value.fetch_market_quotes.return_value = rows + omitted = self.api.handle("/v1/quotes/latest", {}) + empty = self.api.handle("/v1/quotes/latest", {"codes": [""]}) + self.assertEqual(len(omitted["data"]), 220) + self.assertEqual(omitted["meta"]["scope"], "market") + self.assertEqual(omitted["meta"]["source"], "eastmoney:clist") + self.assertEqual(len(empty["data"]), 220) + + def test_explicit_codes_still_use_named_quote_path(self) -> None: + named = [ + { + "ts_code": "600000.SH", + "name": "浦发银行", + "price": 10.2, + "previous_close": 10.0, + } + ] + with patch("datahub.realtime_serve.EastmoneyAdapter") as mocked: + mocked.return_value.fetch_quotes.return_value = named + payload = self.api.handle("/v1/quotes/latest", {"codes": ["600000.SH"]}) + mocked.return_value.fetch_market_quotes.assert_not_called() + self.assertEqual(payload["data"][0]["ts_code"], "600000.SH") + + def test_market_unavailable_stays_source_error(self) -> None: + with patch("datahub.realtime_serve.EastmoneyAdapter") as mocked: + mocked.return_value.fetch_market_quotes.side_effect = AdapterError("too small") + with self.assertRaises(ApiError) as ctx: + self.api.handle("/v1/quotes/latest", {}) + self.assertEqual(ctx.exception.code, "SOURCE_UNAVAILABLE") + + if __name__ == "__main__": unittest.main()