Files
xiaobai-review/backend/data/datahub/bridge.py
T
c8a9376adb fix(HEL-490): 真实装配接通中枢并收编估值晚间复核
把 query/行情钩子绑到内层 TushareClient,图表接受不完整日K窗口;收编现网 HEL-423 未提交的估值复核,避免换版丢掉。

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
2026-09-08 15:03:34 +08:00

543 lines
22 KiB
Python

from __future__ import annotations
import logging
import sys
from typing import Any, Callable
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 (
API_TO_DATASET,
filter_calendar_rows,
filter_stock_rows,
project_fields,
to_native_rows,
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
LOGGER = logging.getLogger("xiaobai.datahub")
ShadowSink = Callable[[dict[str, Any]], None]
def _usable_intraday_points(rows: list[Any]) -> list[dict[str, Any]]:
points: list[dict[str, Any]] = []
for row in rows:
if not isinstance(row, dict):
continue
try:
close = float(row.get("close") or 0)
except (TypeError, ValueError):
close = 0.0
if close <= 0:
continue
point = dict(row)
if "average" not in point and point.get("avg_price") is not None:
point["average"] = point.get("avg_price")
points.append(point)
return points
EMPTY_FAIL_DATASETS = {
"stocks", "daily", "index_daily", "valuation", "moneyflow", "auction",
"limit_events", "sector_daily",
}
def looks_like_heaven(module_name: str, filename: str = "") -> bool:
"""问天调用栈识别(诊断用)。问天按数据集依赖接入,不再整栈强制旧链路。"""
path = filename.replace("\\", "/")
return module_name.startswith("backend.features.heaven") or "/features/heaven/" in path
def caller_is_heaven(depth: int = 24) -> bool:
frame = sys._getframe(1)
for _ in range(depth):
frame = frame.f_back if frame is not None else None
if frame is None:
return False
name = str(frame.f_globals.get("__name__") or "")
filename = str(frame.f_code.co_filename or "")
if looks_like_heaven(name, filename):
return True
return False
class DatahubBridge:
def __init__(
self,
settings: DatahubSettings,
client: DatahubClient,
shadow_sink: ShadowSink | None = None,
heaven_guard: Callable[[], bool] | None = None,
) -> None:
self.settings = settings
self.client = client
self.shadow_sink = shadow_sink
self.heaven_guard = heaven_guard or caller_is_heaven
def dataset_status(self, trade_date: str) -> list[dict[str, Any]] | None:
flags = self.settings.flags("status")
if not flags.read and not flags.shadow:
return None
try:
response = self._require_fresh(self.client.dataset_status(yyyymmdd(trade_date)), "status")
rows = list(response.data or [])
if flags.shadow:
self._emit_shadow(compare_rows("status", [], rows, response.meta))
if flags.read:
return rows
return None
except Exception as exc:
self._log_failure("status", exc)
if flags.shadow:
self._emit_shadow(compare_rows("status", [], [], {}, self._error_text(exc)))
return None
def batches(self, trade_date: str, dataset: str = "") -> list[dict[str, Any]] | None:
flags = self.settings.flags("status")
if not flags.read:
return None
try:
response = self._require_fresh(
self.client.batches(yyyymmdd(trade_date), dataset),
"status",
)
return list(response.data or [])
except Exception as exc:
self._log_failure("status", exc)
return None
def try_intraday(self, code: str) -> dict[str, Any] | None:
flags = self.settings.flags("intraday")
if not flags.read:
return None
try:
response = self.client.intraday_points(code=code)
data = response.data
if not isinstance(data, dict):
raise DatahubError("EMPTY", "datahub intraday payload invalid")
points = _usable_intraday_points(data.get("points") or [])
if not points:
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),
"name": str(data.get("name") or ""),
"code": str(data.get("code") or code),
"trade_date": str(data.get("trade_date") or points[-1].get("date") or ""),
"previous_close": float(data.get("previous_close") or 0),
"points": points,
"source": "datahub",
}
except Exception as exc:
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"},
)
# Charts can use a partial history window; do not discard usable bars
# just because the requested lookback is not fully covered.
self._validate_usable(
dataset,
list(response.data or []),
response,
require_complete=False,
)
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,
params: dict[str, Any] | None,
fields: str,
legacy_query: Callable[..., list[dict[str, Any]]],
) -> list[dict[str, Any]]:
dataset = API_TO_DATASET.get(api_name)
# 问天按实际数据依赖接入:已映射到 hub 的 API 跟随开关;未映射的继续旧链路。
if not dataset:
return legacy_query(api_name, params, fields)
flags = self.settings.flags(dataset)
if not flags.read and not flags.shadow:
return legacy_query(api_name, params, fields)
hub_rows: list[dict[str, Any]] | None = None
hub_meta: dict[str, Any] = {}
hub_error: str | None = None
hub_canonical: list[dict[str, Any]] = []
try:
response = self._fetch_dataset(dataset, params or {}, api_name=api_name)
hub_canonical = self._extract_rows(dataset, response, params or {})
hub_rows = to_native_rows(dataset, hub_canonical)
hub_meta = dict(response.meta)
self._validate_usable(dataset, hub_rows, response)
except Exception as exc:
hub_error = self._error_text(exc)
self._log_failure(dataset, exc)
if flags.shadow:
try:
legacy_rows = legacy_query(api_name, params, fields)
except Exception as exc:
if flags.read and hub_rows is not None and hub_error is None:
self._emit_shadow(
compare_rows(dataset, [], hub_canonical, hub_meta, self._error_text(exc), fields)
)
return project_fields(hub_rows, fields)
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)
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"))
start = yyyymmdd(params.get("start_date") or params.get("from") or date)
end = yyyymmdd(params.get("end_date") or params.get("to") or date)
code = str(params.get("ts_code") or params.get("code") or "").strip()
if dataset == "calendar":
if not start or not end:
raise DatahubError("INVALID_ARGUMENT", "calendar requires start_date and end_date")
return self.client.calendar(start, end)
if dataset == "stocks":
return self._paginate(self.client.stocks, {})
fetchers = {
"daily": self.client.daily_bars,
"index_daily": self.client.index_bars,
"valuation": self.client.valuation,
"moneyflow": self.client.moneyflow,
"auction": self.client.auction,
"limit_events": self.client.limit_events,
"popularity": self.client.popularity,
"dragon_tiger": self.client.dragon_tiger,
"sector_daily": self.client.sectors,
}
fetcher = fetchers[dataset]
query: dict[str, Any] = {}
if code:
query["code"] = code
if date and not (params.get("start_date") or params.get("end_date")):
query["date"] = date
else:
if start:
query["from"] = start
if end:
query["to"] = end
if dataset == "daily":
query["adjust"] = "none"
if dataset == "limit_events":
limit_type = str(params.get("limit_type") or "").strip().upper()
if limit_type:
query["limit_type"] = limit_type
if dataset == "popularity":
if api_name == "ths_hot":
query["source"] = "ths"
elif api_name == "dc_hot":
query["source"] = "dc"
if dataset == "sector_daily":
family = {
"ths_daily": "ths",
"dc_index": "dc",
"sw_daily": "sw",
}.get(api_name, "")
if family:
query["family"] = family
return self._paginate(fetcher, query)
def _paginate(self, fetcher: Callable[..., DatahubResponse], params: dict[str, Any]) -> DatahubResponse:
limit = self.settings.page_limit
offset = 0
rows: list[Any] = []
meta: dict[str, Any] = {}
schema_version = 1
while True:
page = fetcher(**{**params, "limit": limit, "offset": offset})
meta = dict(page.meta)
schema_version = page.schema_version
data = page.data or []
if not isinstance(data, list):
raise DatahubError("INTERNAL", "datahub returned a non-list payload")
rows.extend(data)
if len(data) < limit:
break
offset += limit
if offset > 200_000:
break
return DatahubResponse(data=rows, meta=meta, schema_version=schema_version)
def _extract_rows(
self,
dataset: str,
response: DatahubResponse,
params: dict[str, Any],
) -> list[dict[str, Any]]:
rows = [dict(item) for item in (response.data or [])]
if dataset == "calendar":
return filter_calendar_rows(rows, params)
if dataset == "stocks":
return filter_stock_rows(rows, params)
return rows
def _validate_usable(
self,
dataset: str,
rows: list[dict[str, Any]],
response: DatahubResponse,
require_complete: bool = True,
) -> None:
meta = response.meta or {}
stale_seconds = int(meta.get("staleness_seconds") or 0)
if meta.get("stale") or stale_seconds > self.settings.stale_seconds_max:
raise DatahubError("STALE", f"{dataset} data is stale")
if dataset in EMPTY_FAIL_DATASETS and not rows:
raise DatahubError("EMPTY", f"{dataset} returned no rows")
coverage = meta.get("coverage") if isinstance(meta.get("coverage"), dict) else {}
if require_complete and (meta.get("incomplete") is True or coverage.get("complete") is False):
missing = coverage.get("missing_count")
raise DatahubError("INCOMPLETE", f"{dataset} range is incomplete missing={missing}")
def _require_fresh(self, response: DatahubResponse, dataset: str) -> DatahubResponse:
self._validate_usable(dataset, list(response.data or []) if isinstance(response.data, list) else [], response)
return response
def _emit_shadow(self, report: dict[str, Any]) -> None:
safe = redact_value(report, secrets=self.settings.secrets())
LOGGER.info("datahub shadow %s", safe)
if self.shadow_sink is not None:
self.shadow_sink(report)
def _log_failure(self, dataset: str, exc: Exception) -> None:
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):
text = f"{exc.code}: {exc.message}"
else:
text = str(exc)
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
self._bridge = bridge
# Mixins run as methods on the inner instance (dashboard / indices /
# getattr). Bind hub hooks and query onto that instance so real
# assembly cannot skip 8766.
self._legacy_query = legacy.query
legacy.query = self.query
legacy.try_market_quotes = self.try_market_quotes
legacy.try_quotes = self.try_quotes
legacy.try_index_quotes = self.try_index_quotes
legacy.record_datahub_legacy = self.record_datahub_legacy
def query(
self,
api_name: str,
params: dict[str, Any] | None = None,
fields: str = "",
) -> 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)