Compare commits

..
Author SHA1 Message Date
总工andmultica-agent 8ff537db7b test(HEL-479): scope revision-review batch assertions to A-group boundary for HEL-463 integration
Co-authored-by: multica-agent <github@multica.ai>
2026-09-08 08:25:12 +08:00
总工 260f8464cf Merge commit '605f97e5dffd25132f6a857b7870ab0f023c8f2a' into integrate/hel478-hel463 2026-09-08 08:23:31 +08:00
d175bb65d4 feat(HEL-478): 估值发布后晚间复核并原子追补上游修订
盘后成功发布后继续轻量比对 daily_basic 网站字段,发现修订才走质量门与整组原子切换,避免 17:10 快照落后于晚间上游改写。

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
2026-09-07 21:19:35 +08:00
37 changed files with 190 additions and 3396 deletions
+5 -4
View File
@@ -5,10 +5,11 @@ APP_ENCRYPTION_KEY=
# the system settings; all accounts use the same backend market snapshot.
TUSHARE_TOKEN=your_tushare_token_here
# 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.
# 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.
DATAHUB_BASE_URL=http://127.0.0.1:8766
DATAHUB_TOKEN=
+9 -270
View File
@@ -16,32 +16,11 @@ 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
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",
@@ -112,142 +91,6 @@ class DatahubBridge:
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,
@@ -289,19 +132,12 @@ 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)
result = legacy_query(api_name, params, fields)
if flags.read:
self._record_route(dataset, "legacy", "tushare", hub_error or "")
return result
return legacy_query(api_name, params, fields)
def _fetch_dataset(self, dataset: str, params: dict[str, Any], api_name: str = "") -> DatahubResponse:
date = yyyymmdd(params.get("trade_date") or params.get("date"))
@@ -391,13 +227,7 @@ class DatahubBridge:
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:
def _validate_usable(self, dataset: str, rows: list[dict[str, Any]], response: DatahubResponse) -> 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:
@@ -405,7 +235,7 @@ class DatahubBridge:
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):
if 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}")
@@ -420,12 +250,11 @@ class DatahubBridge:
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()))
LOGGER.warning(
"datahub fallback dataset=%s error=%s",
dataset,
redact_text(self._error_text(exc), self.settings.secrets()),
)
def _error_text(self, exc: Exception) -> str:
if isinstance(exc, DatahubError):
@@ -435,88 +264,10 @@ 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
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,
@@ -524,19 +275,7 @@ class DatahubAwareTushareClient:
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)
return self._bridge.query(api_name, params, fields, self._legacy.query)
def __getattr__(self, name: str) -> Any:
return getattr(self._legacy, name)
-57
View File
@@ -1,57 +0,0 @@
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()
+3 -31
View File
@@ -37,9 +37,7 @@ class DataGateway:
) -> TushareClient:
if dataset_id:
self.policy.assert_allowed(dataset_id, "tushare", usage)
legacy = self.tushare_provider.client()
legacy.realtime_aggregator = self.realtime_observer
return DatahubAwareTushareClient(legacy, self.datahub)
return DatahubAwareTushareClient(self.tushare_provider.client(), self.datahub)
def dataset_status(self, trade_date: str) -> list[dict[str, Any]] | None:
return self.datahub.dataset_status(trade_date)
@@ -47,31 +45,6 @@ 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)
@@ -112,13 +85,12 @@ def build_data_gateway(
policy = DataSourcePolicy.load()
settings = datahub_settings or DatahubSettings.load(credentials=credentials)
datahub_client = DatahubClient(settings)
datahub = DatahubBridge(settings, datahub_client)
return DataGateway(
policy=policy,
quality=DataQualityGate.load(policy),
tushare_provider=TushareProvider(token_supplier),
ifind_provider=IfindProvider(ifind),
chart_data=MarketChartClient(ifind, EastmoneyChartClient(), datahub),
chart_data=MarketChartClient(ifind, EastmoneyChartClient()),
realtime_observer=WebRealtimeAggregator(),
datahub=datahub,
datahub=DatahubBridge(settings, datahub_client),
)
+2 -10
View File
@@ -3,11 +3,7 @@ from __future__ import annotations
from typing import Any
from backend.data.numbers import finite_number as _number
from backend.data.providers.tushare_helpers import (
_display_time,
_prices_equal,
calendar_is_open,
)
from backend.data.providers.tushare_helpers import _display_time, _prices_equal
class DailyMarketMixin:
@@ -21,11 +17,7 @@ class DailyMarketMixin:
trade_date = requested
else:
row = requested_rows[0]
trade_date = (
row["cal_date"]
if calendar_is_open(row.get("is_open"))
else row.get("pretrade_date", requested)
)
trade_date = row["cal_date"] if row.get("is_open") == 1 else row.get("pretrade_date", requested)
resolved_rows = self.query(
"trade_cal",
+12 -148
View File
@@ -16,12 +16,6 @@ from backend.data.providers.tushare_transport import TushareError
class DashboardMixin:
def _now(self) -> datetime:
clock = getattr(self, "clock", None)
if callable(clock):
return clock()
return datetime.now().astimezone()
def dashboard(self, requested_date: str) -> dict[str, Any]:
trade_date, previous_trade_date = self.resolve_trade_context(requested_date)
if self.should_use_realtime(requested_date, trade_date):
@@ -32,12 +26,11 @@ class DashboardMixin:
)
daily = self._load_daily(trade_date)
now = self._now()
if (
not daily
and requested_date == now.strftime("%Y%m%d")
and requested_date == datetime.now().astimezone().strftime("%Y%m%d")
and trade_date == requested_date
and now.time().replace(tzinfo=None) >= dt_time(9, 15)
and datetime.now().astimezone().time().replace(tzinfo=None) >= dt_time(9, 15)
):
return self._realtime_dashboard(
requested_date,
@@ -105,14 +98,15 @@ class DashboardMixin:
}
return apply_sentiment_to_dashboard(dashboard)
def should_use_realtime(self, requested_date: str, trade_date: str) -> bool:
"""Use live quotes for today's open session until official daily settles."""
now = self._now()
@staticmethod
def should_use_realtime(requested_date: str, trade_date: str) -> bool:
"""Use rt_k for today's open market until end-of-day datasets settle."""
now = datetime.now().astimezone()
today = now.strftime("%Y%m%d")
return (
requested_date == today
and trade_date == today
and dt_time(9, 15) <= now.time().replace(tzinfo=None) < dt_time(15, 5)
and dt_time(9, 15) <= now.time().replace(tzinfo=None) < dt_time(16, 30)
)
def _realtime_dashboard(
@@ -128,7 +122,7 @@ class DashboardMixin:
)
if not codes:
raise TushareError("No active stock codes available for rt_k")
quotes, quote_source = self._load_realtime_quotes(codes, trade_date)
quotes = self.query("rt_k", {"ts_code": codes})
if not quotes:
raise TushareError(f"No realtime data returned for {trade_date}")
@@ -184,35 +178,14 @@ class DashboardMixin:
)
sectors = _build_sectors(limits)
previous_sectors = _build_sectors(previous_limits)
now = self._now()
now = datetime.now().astimezone()
market_status = _realtime_market_status(now.time().replace(tzinfo=None))
if quote_source == "datahub":
notice = (
"盘中行情由数据中枢统一提供;涨停原因、封板时间和开板次数以盘后榜单校正为准。"
)
source_name = "datahub"
elif 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": source_name,
"quote_source": quote_source,
"source": "tushare",
"mode": "realtime",
"realtime": True,
"market_status": market_status,
@@ -220,8 +193,7 @@ class DashboardMixin:
"auto_refresh": False,
"quote_count": len(daily),
"updated_at": now.isoformat(timespec="seconds"),
"notice": notice,
"indices": self._free_realtime_indices() if quote_source != "tushare_rt_k" else [],
"notice": "盘中行情由 Tushare rt_k 实时计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。",
},
"overview": _build_overview(daily, up_rows, down_rows, broken_rows),
"limits": limits,
@@ -235,89 +207,6 @@ 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]:
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:
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"
)
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,
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]]:
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:
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 []
def _load_realtime_reference(
self,
trade_date: str,
@@ -345,7 +234,7 @@ class DashboardMixin:
{"trade_date": previous_trade_date},
"ts_code,trade_date,total_share,float_share,free_share,total_mv,circ_mv",
)
if not basic_rows:
if not basic_rows or not price_limits:
raise TushareError(f"Realtime reference data is incomplete for {trade_date}")
result = {
"basic_rows": basic_rows,
@@ -719,31 +608,6 @@ 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):
-11
View File
@@ -6,17 +6,6 @@ from typing import Any
from backend.data.numbers import finite_number as _number
def calendar_is_open(value: Any) -> bool:
if value in (True, 1, "1", "Y", "y"):
return True
if value in (False, 0, "0", "N", "n", None, ""):
return False
try:
return int(value) == 1
except (TypeError, ValueError):
return False
def _text(value: Any) -> str:
if isinstance(value, (list, tuple, set)):
return "".join(str(item).strip() for item in value if str(item).strip())
-130
View File
@@ -59,87 +59,6 @@ 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:
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:
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)
index_names = {
"000001.SH": "上证指数",
@@ -197,52 +116,3 @@ 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,
},
}
-324
View File
@@ -19,20 +19,8 @@ class RealtimeAggregateError(RuntimeError):
EASTMONEY_INDEX_URL = "https://push2.eastmoney.com/api/qt/ulist.np/get"
EASTMONEY_STOCK_URL = "https://push2.eastmoney.com/api/qt/stock/get"
EASTMONEY_STOCK_FIELDS = "f43,f44,f45,f46,f47,f48,f57,f58,f60,f86,f168"
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 = (
@@ -146,181 +134,6 @@ 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_stock_quote(self, code: str, expected_date: str = "") -> dict[str, Any]:
symbol, _secid, ts_code = _a_share_identity(code)
raw, _cache_age = self._get_text(
f"{TENCENT_QUOTE_URL}{symbol}",
referer="https://gu.qq.com/",
encoding="gb18030",
)
quote = next(
(
item
for line in raw.splitlines()
if (item := _parse_tencent_stock_quote(line))
),
None,
)
if not quote:
raise RealtimeAggregateError(f"Tencent stock quote unavailable for {ts_code}")
return _require_quote_date(quote, expected_date)
def eastmoney_stock_quote(self, code: str, expected_date: str = "") -> dict[str, Any]:
_symbol, secid, ts_code = _a_share_identity(code)
payload = self._get_json(
EASTMONEY_STOCK_URL,
{
"secid": secid,
"invt": "2",
"fltt": "2",
"fields": EASTMONEY_STOCK_FIELDS,
},
referer="https://quote.eastmoney.com/",
)
quote = _normalize_eastmoney_stock_quote(payload.get("data") or {}, ts_code)
if not quote:
raise RealtimeAggregateError(f"Eastmoney stock quote unavailable for {ts_code}")
return _require_quote_date(quote, expected_date)
def tencent_indices(self) -> list[dict[str, Any]]:
raw, cache_age = self._get_text(
TENCENT_INDEX_URL,
@@ -584,143 +397,6 @@ 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 _a_share_identity(code: str) -> tuple[str, str, str]:
raw = str(code or "").strip().upper()
symbol = raw.split(".")[0]
if not symbol.isdigit() or len(symbol) != 6:
raise RealtimeAggregateError("Invalid stock code")
if raw.endswith(".SH") or symbol.startswith(("5", "6", "9")):
return f"sh{symbol}", f"1.{symbol}", f"{symbol}.SH"
if raw.endswith(".BJ") or symbol.startswith(("4", "8")):
return f"bj{symbol}", f"0.{symbol}", f"{symbol}.BJ"
return f"sz{symbol}", f"0.{symbol}", f"{symbol}.SZ"
def _require_quote_date(quote: dict[str, Any], expected_date: str) -> dict[str, Any]:
want = str(expected_date or "").replace("-", "")
got = str(quote.get("quote_date") or "")
if want and got != want:
raise RealtimeAggregateError(f"quote date {got or 'empty'} is not {want}")
return quote
def _normalize_eastmoney_stock_quote(
row: dict[str, Any], ts_code: str
) -> dict[str, Any] | None:
close = _number(row.get("f43"))
previous_close = _number(row.get("f60"))
if close <= 0 or previous_close <= 0:
return None
epoch = int(_number(row.get("f86")))
quote_date = ""
if epoch > 0:
quote_date = datetime.fromtimestamp(epoch).astimezone().strftime("%Y%m%d")
return {
"ts_code": ts_code,
"name": row.get("f58") or ts_code.split(".")[0],
"pre_close": previous_close,
"open": _number(row.get("f46")),
"high": _number(row.get("f44")),
"low": _number(row.get("f45")),
"close": close,
"vol": _number(row.get("f47")) * 100,
"amount": _number(row.get("f48")),
"num": 0,
"quote_date": quote_date,
"quote_time_epoch": epoch,
"turnover_rate": _number(row.get("f168")),
"source": "eastmoney_stock",
}
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股)"):
+15 -175
View File
@@ -2,7 +2,6 @@ from __future__ import annotations
import http.client
import json
import logging
import re
import time
import urllib.error
@@ -16,15 +15,12 @@ from typing import Any, ClassVar
from backend.bootstrap.config import tushare_code as _stock_market_code
from backend.data.providers.ifind_client import IfindError, IfindHttpClient
LOGGER = logging.getLogger("xiaobai.charts")
class ChartDataError(RuntimeError):
pass
TRENDS_URL = "https://push2delay.eastmoney.com/api/qt/stock/trends2/get"
HIS_TRENDS_URL = "https://push2his.eastmoney.com/api/qt/stock/trends2/get"
BOARD_LIST_URL = "https://push2delay.eastmoney.com/api/qt/clist/get"
BROWSER_USER_AGENT = (
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
@@ -41,23 +37,14 @@ INDEX_SECIDS = {
class MarketChartClient:
"""Prefer iFinD for display charts and retain Eastmoney as a last resort."""
def __init__(
self,
ifind: IfindHttpClient,
fallback: "EastmoneyChartClient",
datahub: Any = None,
) -> None:
def __init__(self, ifind: IfindHttpClient, fallback: "EastmoneyChartClient") -> None:
self.ifind = ifind
self.fallback = fallback
self.datahub = datahub
def stock_intraday(self, code: str) -> dict[str, Any]:
normalized = str(code or "").strip()
if not re.fullmatch(r"\d{6}", normalized):
raise ChartDataError("Invalid stock code")
hub_chart = self._datahub_intraday(normalized)
if hub_chart is not None:
return hub_chart
ifind_code = _stock_market_code(normalized)
try:
return self._ifind_intraday(ifind_code, "stock", normalized)
@@ -68,18 +55,12 @@ 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]]:
@@ -92,135 +73,11 @@ class MarketChartClient:
normalized = str(identifier or "").strip().upper()
if normalized not in INDEX_SECIDS:
raise ChartDataError("Unsupported index")
hub_chart = self._datahub_intraday(normalized)
if hub_chart is not None:
return hub_chart
try:
return self._ifind_intraday(normalized, "index", normalized)
except (IfindError, ChartDataError):
return self.fallback.index_intraday(normalized)
def _datahub_intraday(self, code: str) -> dict[str, Any] | None:
if self.datahub is None:
return None
try:
chart = self.datahub.try_intraday(code)
except Exception as exc:
LOGGER.warning("datahub intraday unexpected error: %s", exc)
return None
if not chart:
return None
points = list(chart.get("points") or [])
if not points:
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:
@@ -448,29 +305,21 @@ class EastmoneyChartClient:
if cached is not None:
return cached
params = {
"secid": secid,
"fields1": "f1,f2,f3,f4,f5,f6,f7,f8,f9,f10,f11,f12,f13",
"fields2": "f51,f52,f53,f54,f55,f56,f57,f58",
"iscr": "0",
}
last_error: Exception | None = None
data: dict[str, Any] = {}
points: list[dict[str, Any]] = []
for url, ndays in ((TRENDS_URL, "1"), (TRENDS_URL, "5"), (HIS_TRENDS_URL, "5")):
request_params = {**params, "ndays": ndays}
try:
payload = self._request_json(url, request_params, "https://quote.eastmoney.com/")
except ChartDataError as exc:
last_error = exc
continue
data = payload.get("data") or {}
parsed = [point for raw in data.get("trends") or [] if (point := _parse_trend(raw))]
points = _latest_session(parsed)
if points:
break
payload = self._request_json(
TRENDS_URL,
{
"secid": secid,
"fields1": "f1,f2,f3,f4,f5,f6,f7,f8,f9,f10,f11,f12,f13",
"fields2": "f51,f52,f53,f54,f55,f56,f57,f58",
"iscr": "0",
"ndays": "1",
},
"https://quote.eastmoney.com/",
)
data = payload.get("data") or {}
points = [point for raw in data.get("trends") or [] if (point := _parse_trend(raw))]
if not points:
raise ChartDataError("No intraday chart data returned") from last_error
raise ChartDataError("No intraday chart data returned")
result = {
"entity_type": entity_type,
@@ -584,15 +433,6 @@ class EastmoneyChartClient:
raise ChartDataError("Intraday chart request failed") from last_error
def _latest_session(points: list[dict[str, Any]]) -> list[dict[str, Any]]:
if not points:
return []
latest = max(str(point.get("date") or "") for point in points)
if not latest:
return points
return [point for point in points if str(point.get("date") or "") == latest]
def _parse_trend(raw: Any) -> dict[str, Any] | None:
fields = str(raw or "").split(",")
if len(fields) < 8 or " " not in fields[0]:
+21 -215
View File
@@ -15,7 +15,6 @@ from backend.bootstrap.config import (
)
from backend.data.providers.ifind_client import IfindError
from backend.data.providers.tushare_client import TushareClient, TushareError
from backend.data.realtime import RealtimeAggregateError
from backend.features.market.backfill_history import (
DEFAULT_RECENT_TRADING_DAYS,
MAX_RANGE_TRADING_DAYS,
@@ -43,7 +42,6 @@ SEARCH_TYPE_LABELS = {
"theme": "题材",
"index": "指数",
}
TODAY_DAILY_UNAVAILABLE_NOTICE = "今日日K暂不可用,仍显示最近收盘K线。"
THS_SEARCH_TYPES = {
"I": ("sector", "行业板块"),
"R": ("sector", "地域板块"),
@@ -65,37 +63,11 @@ class MarketServiceMixin:
if gateway is not None:
return gateway.tushare()
# Compatibility for isolated legacy unit-test service stubs.
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)
if callable(clock):
return clock()
return datetime.now().astimezone()
def _is_requested_open_session(self, requested_date: str) -> bool:
now = self._now()
if requested_date != now.strftime("%Y%m%d"):
return False
if now.time().replace(tzinfo=None) < dt_time(9, 15):
return False
client = self._tushare_client() if self.configured else None
resolve = getattr(client, "resolve_trade_context", None) if client else None
if resolve is None:
return now.weekday() < 5
try:
trade_date, _ = resolve(requested_date)
except Exception:
return now.weekday() < 5
return str(trade_date or "") == requested_date
return TushareClient(self.token)
def get_dashboard(self, trade_date: str, force: bool = False) -> dict[str, Any]:
normalized_date = normalize_date(trade_date)
now = self._now()
now = datetime.now().astimezone()
if (
normalized_date == now.strftime("%Y%m%d")
and now.time().replace(tzinfo=None) < datetime.strptime("09:15", "%H:%M").time()
@@ -202,14 +174,14 @@ class MarketServiceMixin:
def _should_retry_incomplete_snapshot(
self, snapshot: dict[str, Any], requested_date: str
) -> bool:
if requested_date != self._now().strftime("%Y%m%d"):
if requested_date != date.today().strftime("%Y%m%d"):
return False
meta = snapshot.get("meta") or {}
actual = str(meta.get("trade_date") or "").replace("-", "")
stale_carry = bool(meta.get("carried_forward") or actual != requested_date)
if stale_carry and self._is_requested_open_session(requested_date):
return True
incomplete = meta.get("limit_data_source") == "derived" or stale_carry
incomplete = (
meta.get("limit_data_source") == "derived"
or bool(meta.get("carried_forward"))
or str(meta.get("trade_date") or "").replace("-", "") != requested_date
)
return incomplete and self._snapshot_age_seconds(meta) >= 60
def _annotate_data_status(self, dashboard: dict[str, Any]) -> dict[str, Any]:
@@ -227,9 +199,6 @@ class MarketServiceMixin:
else:
meta["data_status"] = "preparing"
meta["display_notice"] = self._preparing_display_notice(actual, requested)
elif meta.get("realtime"):
meta["data_status"] = "intraday"
meta.setdefault("display_notice", "")
else:
meta["data_status"] = "official"
meta.setdefault("display_notice", "")
@@ -256,9 +225,9 @@ class MarketServiceMixin:
normalized_date: str,
snapshot: dict[str, Any],
) -> bool:
if not self.configured or normalized_date != self._now().strftime("%Y%m%d"):
if not self.configured or normalized_date != date.today().strftime("%Y%m%d"):
return False
now = self._now()
now = datetime.now().astimezone()
local_time = now.time().replace(tzinfo=None)
realtime_start = datetime.strptime("09:15", "%H:%M").time()
morning_end = datetime.strptime("11:35", "%H:%M").time()
@@ -295,10 +264,7 @@ 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(
@@ -310,12 +276,6 @@ class MarketServiceMixin:
actual_date = normalize_date(
str(dashboard.get("meta", {}).get("trade_date") or normalized_date)
)
if actual_date != normalized_date and self._is_requested_open_session(
normalized_date
):
raise TushareError(
f"Intraday dashboard resolved {actual_date} instead of {normalized_date}"
)
self.database.save_snapshot(actual_date, source, dashboard)
if actual_date != normalized_date:
dashboard.setdefault("meta", {}).update(
@@ -337,30 +297,6 @@ class MarketServiceMixin:
)
return self._apply_reason_overrides(self._with_storage(dashboard, cached=False))
except TushareError as exc:
if self._is_requested_open_session(normalized_date):
existing = self.database.get_snapshot(normalized_date)
existing_date = str(
((existing or {}).get("meta") or {}).get("trade_date") or ""
).replace("-", "")
if existing and existing_date == normalized_date:
kept = copy.deepcopy(existing)
kept.setdefault("meta", {}).update(
{
"requested_date": self._display_compact_date(normalized_date),
}
)
self.database.finish_sync(
sync_id,
"fallback",
self._record_count(kept),
str(exc),
"tushare",
)
return self._apply_reason_overrides(
self._with_storage(kept, cached=True)
)
self.database.finish_sync(sync_id, "failed", message=str(exc))
raise ValueError("当天盘中行情暂时不可用,请稍后重试。") from exc
fallback = self.database.get_latest_real_snapshot(normalized_date)
if fallback:
actual = str((fallback.get("meta") or {}).get("trade_date") or "")
@@ -816,27 +752,26 @@ class MarketServiceMixin:
"trade_date": f"{actual_date[:4]}-{actual_date[4:6]}-{actual_date[6:]}",
}
today = now.strftime("%Y%m%d")
latest_bar = (result.get("prices") or [{}])[-1] if result.get("prices") else {}
official_today = (
actual_date == today and not bool(latest_bar.get("realtime"))
)
after_close = now.time().replace(tzinfo=None) >= dt_time(15, 0)
should_merge = (
requested_date == today
and actual_date <= today
and now.weekday() < 5
and now.time().replace(tzinfo=None) >= dt_time(9, 30)
and not (official_today and after_close)
)
if should_merge:
quote = self._resolve_today_daily_quote(code, today, result)
quote = self._ifind_realtime_stock_quote(code)
if quote and self._valid_realtime_stock_quote(quote, today):
self._merge_realtime_stock_detail(result, quote, requested_date)
elif actual_date < today:
result["meta"] = {
**(result.get("meta") or {}),
"notice": TODAY_DAILY_UNAVAILABLE_NOTICE,
}
elif self.configured and actual_date < today:
client = self._tushare_client()
try:
resolved_date, _ = client.resolve_trade_context(requested_date)
if resolved_date == today:
quote = client.realtime_stock_quote(tushare_code(code), requested_date)
if self._valid_realtime_stock_quote(quote, today):
self._merge_realtime_stock_detail(result, quote, requested_date)
except TushareError:
pass
return self._enrich_stock_detail(result)
@staticmethod
@@ -951,134 +886,6 @@ class MarketServiceMixin:
"quote_time": str(row.get("time") or ""),
}
def _resolve_today_daily_quote(
self, code: str, today: str, payload: dict[str, Any]
) -> dict[str, Any] | None:
quote = self._ifind_realtime_stock_quote(code)
if quote and self._valid_realtime_stock_quote(quote, today):
return quote
if self.configured:
try:
client = self._tushare_client()
resolve = getattr(client, "resolve_trade_context", None)
resolved = today
if callable(resolve):
resolved, _ = resolve(today)
if str(resolved or "") == today:
quote = client.realtime_stock_quote(tushare_code(code), today)
if self._valid_realtime_stock_quote(quote, today):
return quote
except TushareError:
pass
quote = self._free_realtime_stock_quote(code, today)
if quote and self._valid_realtime_stock_quote(quote, today):
return quote
return self._intraday_realtime_stock_quote(code, today, payload)
def _free_realtime_stock_quote(self, code: str, today: str) -> dict[str, Any] | None:
aggregator = getattr(self, "realtime_aggregator", None)
if aggregator is None:
return None
ts_code = tushare_code(code)
for loader in (
getattr(aggregator, "tencent_stock_quote", None),
getattr(aggregator, "eastmoney_stock_quote", None),
):
if not callable(loader):
continue
try:
row = loader(ts_code, expected_date=today)
except (RealtimeAggregateError, Exception):
continue
quote = self._quote_from_free_row(code, today, row)
if quote:
return quote
return None
def _quote_from_free_row(
self, code: str, today: str, row: dict[str, Any]
) -> dict[str, Any] | None:
price = float(row.get("close") or 0)
previous_close = float(row.get("pre_close") or 0)
if price <= 0 or previous_close <= 0:
return None
try:
name, sector = self._stock_identity(code, today)
except Exception:
name, sector = "--", "其他"
epoch = int(row.get("quote_time_epoch") or 0)
if epoch > 0:
quote_time = datetime.fromtimestamp(epoch).astimezone().isoformat(timespec="seconds")
else:
quote_date = str(row.get("quote_date") or today)
quote_time = f"{quote_date[:4]}-{quote_date[4:6]}-{quote_date[6:]}"
return {
"name": str(row.get("name") or name or "--"),
"sector": sector,
"price": price,
"open": float(row.get("open") or 0),
"high": float(row.get("high") or 0),
"low": float(row.get("low") or 0),
"change": round((price / previous_close - 1) * 100, 4),
"volume": float(row.get("vol") or 0),
"amount_billion": float(row.get("amount") or 0) / 100_000_000,
"turnover_rate": float(row.get("turnover_rate") or 0),
"quote_time": quote_time,
}
def _intraday_realtime_stock_quote(
self, code: str, today: str, payload: dict[str, Any]
) -> dict[str, Any] | None:
chart_data = getattr(self, "chart_data", None)
if chart_data is None:
return None
try:
chart = chart_data.stock_intraday(code)
except (AttributeError, ChartDataError, Exception):
return None
points = [
point
for point in list(chart.get("points") or [])
if str(point.get("date") or "").replace("-", "") == today
]
if not points:
return None
opens = [float(point.get("open") or 0) for point in points if float(point.get("open") or 0) > 0]
highs = [float(point.get("high") or 0) for point in points if float(point.get("high") or 0) > 0]
lows = [float(point.get("low") or 0) for point in points if float(point.get("low") or 0) > 0]
closes = [float(point.get("close") or 0) for point in points if float(point.get("close") or 0) > 0]
if not opens or not highs or not lows or not closes:
return None
price = closes[-1]
previous_close = float(chart.get("previous_close") or 0)
if previous_close <= 0:
history = list(payload.get("prices") or [])
previous_close = float((history[-1] if history else {}).get("close") or 0)
if previous_close <= 0:
return None
volume = sum(float(point.get("volume") or 0) for point in points)
amount = sum(float(point.get("amount") or 0) for point in points)
if volume <= 0 and amount <= 0:
return None
try:
name, sector = self._stock_identity(code, today)
except Exception:
name, sector = "--", "其他"
return {
"name": name,
"sector": sector,
"price": price,
"open": opens[0],
"high": max(highs),
"low": min(lows),
"change": round((price / previous_close - 1) * 100, 4),
"volume": volume,
"volume_unit": "lots",
"amount_billion": amount / 100_000_000,
"turnover_rate": 0.0,
"quote_time": str(points[-1].get("date") or today),
}
@staticmethod
def _merge_realtime_stock_detail(
payload: dict[str, Any], quote: dict[str, Any], trade_date: str
@@ -1117,7 +924,6 @@ class MarketServiceMixin:
**(payload.get("meta") or {}),
"trade_date": display_date,
"realtime": True,
"notice": "",
"updated_at": datetime.now().astimezone().isoformat(timespec="seconds"),
}
-17
View File
@@ -130,7 +130,6 @@ class SystemServiceMixin:
),
**self.database.status(),
"jobs": self.jobs.repository.recent(12),
"datahub": self._datahub_status(),
},
"llm": {
"primary_configured": self._profile_configured(platform["primary"]),
@@ -146,22 +145,6 @@ 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()
-2
View File
@@ -41,8 +41,6 @@ def official_catchup_due(today: str, snapshot: dict[str, object]) -> bool:
actual == today
and meta.get("limit_data_source") != "derived"
and not meta.get("carried_forward")
and not meta.get("realtime")
and meta.get("mode") != "realtime"
):
return False
return True
-16
View File
@@ -13,22 +13,6 @@ 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:
+3 -6
View File
@@ -12,12 +12,9 @@ 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`: 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.
- `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.
- `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,
+46 -46
View File
@@ -222,12 +222,12 @@
{
"provider": "eastmoney",
"path": "backend/data/realtime.py",
"runtime_role": "isolated realtime observation and intraday dashboard fallback"
"runtime_role": "isolated realtime observation"
},
{
"provider": "tencent",
"path": "backend/data/realtime.py",
"runtime_role": "index observation and intraday quote fallback"
"runtime_role": "index observation fallback"
}
],
"provider_domains": [
@@ -483,8 +483,8 @@
},
{
"path": "frontend/index.html",
"bytes": 48447,
"lines": 665
"bytes": 48254,
"lines": 664
},
{
"path": "backend/features/screener/catalog.py",
@@ -496,11 +496,6 @@
"bytes": 35247,
"lines": 2416
},
{
"path": "backend/data/providers/tushare_dashboard.py",
"bytes": 33603,
"lines": 784
},
{
"path": "database.py",
"bytes": 32073,
@@ -511,6 +506,11 @@
"bytes": 31756,
"lines": 562
},
{
"path": "backend/data/providers/tushare_dashboard.py",
"bytes": 28234,
"lines": 648
},
{
"path": "backend/data/providers/tushare_industries.py",
"bytes": 26540,
@@ -533,8 +533,8 @@
},
{
"path": "frontend/pages/market/preview.js",
"bytes": 18339,
"lines": 450
"bytes": 18178,
"lines": 446
},
{
"path": "backend/features/heaven/trend.py",
@@ -551,16 +551,6 @@
"bytes": 15311,
"lines": 387
},
{
"path": "frontend/shared/admin.js",
"bytes": 15235,
"lines": 289
},
{
"path": "frontend/shared/dashboard.js",
"bytes": 15063,
"lines": 321
},
{
"path": "frontend/pages/pools/page.html",
"bytes": 14942,
@@ -571,6 +561,16 @@
"bytes": 14743,
"lines": 342
},
{
"path": "frontend/shared/dashboard.js",
"bytes": 14740,
"lines": 316
},
{
"path": "frontend/shared/admin.js",
"bytes": 14410,
"lines": 268
},
{
"path": "backend/features/heaven/market_context.py",
"bytes": 13681,
@@ -581,20 +581,15 @@
"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/data/providers/tushare_indices.py",
"bytes": 10956,
"lines": 248
"path": "backend/features/system/service.py",
"bytes": 12392,
"lines": 254
},
{
"path": "backend/features/market/insights_auction.py",
@@ -643,8 +638,8 @@
},
{
"path": "backend/data/providers/tushare_daily.py",
"bytes": 6949,
"lines": 168
"bytes": 6837,
"lines": 160
},
{
"path": "backend/application.py",
@@ -681,16 +676,21 @@
"bytes": 6092,
"lines": 138
},
{
"path": "frontend/pages/market/stock-detail.js",
"bytes": 6041,
"lines": 134
},
{
"path": "frontend/pages/dragon-tiger/page.html",
"bytes": 5754,
"lines": 85
},
{
"path": "frontend/pages/market/stock-detail.js",
"bytes": 5690,
"lines": 124
},
{
"path": "backend/data/providers/tushare_indices.py",
"bytes": 5451,
"lines": 118
},
{
"path": "frontend/pages.config.js",
"bytes": 5385,
@@ -786,11 +786,6 @@
"bytes": 2514,
"lines": 63
},
{
"path": "backend/data/providers/tushare_helpers.py",
"bytes": 2360,
"lines": 75
},
{
"path": "backend/jobs/service.py",
"bytes": 2337,
@@ -816,6 +811,11 @@
"bytes": 2165,
"lines": 35
},
{
"path": "backend/data/providers/tushare_helpers.py",
"bytes": 2083,
"lines": 64
},
{
"path": "frontend/pages/market/breadth.js",
"bytes": 2071,
@@ -826,16 +826,16 @@
"bytes": 1919,
"lines": 45
},
{
"path": "backend/jobs/refresh.py",
"bytes": 1808,
"lines": 48
},
{
"path": "backend/features/system/routes.py",
"bytes": 1791,
"lines": 46
},
{
"path": "backend/jobs/refresh.py",
"bytes": 1728,
"lines": 46
},
{
"path": "backend/features/alerts/routes.py",
"bytes": 1687,
+15 -15
View File
@@ -6,20 +6,20 @@
"page_limit": 5000,
"stale_seconds_max": 86400,
"datasets": {
"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 }
"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 }
}
}
+2 -2
View File
@@ -213,12 +213,12 @@
{
"provider": "eastmoney",
"path": "realtime_aggregator.py",
"runtime_role": "isolated realtime observation and intraday dashboard fallback"
"runtime_role": "isolated realtime observation"
},
{
"provider": "tencent",
"path": "realtime_aggregator.py",
"runtime_role": "index observation and intraday quote fallback"
"runtime_role": "index observation fallback"
}
],
"llm_entrypoints": [
-1
View File
@@ -611,7 +611,6 @@
<label class="form-field"><span>iFinD Refresh Token</span><input id="systemIfindTokenInput" type="password" autocomplete="off" maxlength="2048" placeholder="留空保留现有 Token"></label>
<label class="switch-control"><input id="systemBackgroundRefresh" type="checkbox"><span>启用交易时段后台刷新</span></label>
<p class="form-hint">所有用户读取同一份后台快照,页面不会随后台任务自动重绘。</p>
<div id="datahubRouteStatus" class="admin-refresh-status" data-tone="idle" role="status" aria-live="polite"><i data-lucide="database"></i><span>数据中枢线路待检查</span></div>
<div id="adminRefreshStatus" class="admin-refresh-status" data-tone="idle" role="status" aria-live="polite"><i data-lucide="circle-dot"></i><span>尚未手动刷新</span></div>
<div class="dialog-actions admin-inline-actions"><button id="adminRefreshButton" class="button" type="button"><i data-lucide="refresh-cw"></i>立即后台刷新</button><button class="button primary" type="submit">保存行情配置</button></div>
</form>
-10
View File
@@ -5219,15 +5219,6 @@
return '<span class="m-sys-dot' + (ok ? " m-sys-dot--ok" : "") + '"></span>';
}
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();
@@ -5246,7 +5237,6 @@
'<div class="m-sys-status-item"><span>iFinD</span><span>' + statusDot(ifind.configured) + (ifind.configured ? " 已配置" : " 未配置") + "</span></div>" +
'<div class="m-sys-status-item"><span>行情快照</span><strong>' + number(data.snapshot_dates) + " 个交易日</strong></div>" +
'<div class="m-sys-status-item"><span>后台刷新</span><span>' + statusDot(data.background_refresh_enabled) + (data.background_refresh_enabled ? " 已启用" : " 已暂停") + "</span></div>" +
'<div class="m-sys-status-item"><span>数据中枢</span><span>' + statusDot(Boolean((data.datahub || {}).configured) && !((data.datahub || {}).fallback_count)) + datahubStatusText(data.datahub || {}) + "</span></div>" +
"</div></div>" +
'<div class="m-card m-sys-section"><strong>数据源密钥</strong>' +
formFieldHtml("Tushare Token", '<input id="m-sys-token" type="password" autocomplete="off" minlength="20" placeholder="留空则保留现有 Token">', false) +
+1 -5
View File
@@ -367,11 +367,7 @@ function selectStockPreviewChart(chart) {
}
} else if ((payload.prices || []).length) {
setText("stockPreviewDate", payload.meta?.trade_date || "最新行情");
const notice = String(payload.meta?.notice || "").trim();
setText(
"stockPreviewSource",
notice ? `日 K 行情 · ${payload.prices.length} 个交易日 · ${notice}` : `日 K 行情 · ${payload.prices.length} 个交易日`,
);
setText("stockPreviewSource", `日 K 行情 · ${payload.prices.length} 个交易日`);
drawDailyPreviewChart(payload.prices);
} else {
setText("stockPreviewDate", payload.meta?.trade_date || "最新行情");
+2 -12
View File
@@ -52,11 +52,7 @@ async function openStock(code, fallback = null) {
renderStockNotes(payload.notes || []);
updateWatchButton();
if (state.stockDetailChartMode === "daily") {
const notice = String(payload.meta?.notice || "").trim();
setText(
"chartSource",
notice ? `日 K 行情 · ${payload.prices.length} 个交易日 · ${notice}` : `日 K 行情 · ${payload.prices.length} 个交易日`,
);
setText("chartSource", `日 K 行情 · ${payload.prices.length} 个交易日`);
requestAnimationFrame(() => drawPriceChart(payload.prices || []));
}
} catch (error) {
@@ -73,13 +69,7 @@ async function selectStockDetailChart(mode) {
syncDetailChartButtons("stock", selected);
if (selected === "daily") {
const prices = state.stockDetail?.prices || [];
const notice = String(state.stockDetail?.meta?.notice || "").trim();
setText(
"chartSource",
prices.length
? (notice ? `日 K 行情 · ${prices.length} 个交易日 · ${notice}` : `日 K 行情 · ${prices.length} 个交易日`)
: "正在加载行情",
);
setText("chartSource", prices.length ? `日 K 行情 · ${prices.length} 个交易日` : "正在加载行情");
if (prices.length) requestAnimationFrame(() => drawPriceChart(prices));
else clearPriceChart("正在加载日 K 数据");
return;
-21
View File
@@ -44,7 +44,6 @@ 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);
@@ -56,26 +55,6 @@ 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;
-5
View File
@@ -67,11 +67,6 @@ async function startAdminRefresh() {
const actualCompact = actualDate.replaceAll("-", "");
const updated = formatTimestamp(meta.updated_at);
const freshness = dashboardFreshnessMessage(meta);
if (meta.realtime && actualCompact === requestedCompact && !meta.carried_forward) {
setAdminRefreshStatus("success", `刷新成功:已获取 ${actualDate} 的盘中行情,更新时间 ${updated}`, "circle-check");
showToast(`刷新成功:已获取 ${actualDate} 的盘中行情`);
return;
}
if (freshness || actualCompact !== requestedCompact || meta.carried_forward || meta.limit_data_source === "derived") {
setAdminRefreshStatus("warning", freshness || `部分正式数据尚未到齐,当前展示 ${actualDate || "最近可用数据"}`, "triangle-alert");
setStatus(freshness || "部分正式数据尚未到齐,当前展示最近可用数据");
+15 -244
View File
@@ -3,8 +3,7 @@ from __future__ import annotations
import copy
import threading
import unittest
from datetime import date, datetime, timedelta, timezone, time as dt_time
from unittest.mock import patch
from datetime import date, datetime, timedelta, timezone
from pathlib import Path
from backend.features.market.service import MarketServiceMixin
@@ -106,84 +105,18 @@ class FakeDerivedClient:
}
SHANGHAI = timezone(timedelta(hours=8))
TRADE_DAY = date(2026, 9, 8)
def at_clock(hour: int, minute: int, day: date = TRADE_DAY) -> datetime:
return datetime(day.year, day.month, day.day, hour, minute, tzinfo=SHANGHAI)
class FakeMissingDailyClient:
def __init__(self, open_today: bool = True):
self.open_today = open_today
def dashboard(self, trade_date: str):
raise TushareError(f"No daily data returned for {trade_date}")
def resolve_trade_context(self, requested: str):
if self.open_today:
return requested, "20260907"
return "20260907", "20260904"
class FakeRealtimeTodayClient:
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",
"market_status": "trading",
"notice": "盘中行情由 Tushare rt_k 实时计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。",
"updated_at": datetime.now().astimezone().isoformat(timespec="seconds"),
},
"overview": {"limit_up_count": 15},
"limits": [{"code": "000001"}],
"broken": [],
"down_limits": [],
"yesterday_limits": [],
}
def resolve_trade_context(self, requested: str):
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):
def __init__(self, client, latest=None):
self.configured = True
self.sync_lock = threading.Lock()
self.database = FakeSyncDatabase(latest)
self._client = client
self.current_user_id = 1
self.clock = clock
def _tushare_client(self):
return self._client
@@ -209,161 +142,23 @@ class DashboardFreshnessTests(unittest.TestCase):
self.assertEqual(harness.database.finished[0][0][1], "success")
self.assertEqual(verified_dashboard_result(payload), payload)
def test_intraday_refresh_keeps_today_and_does_not_fall_back_to_yesterday(self):
today = TRADE_DAY.strftime("%Y%m%d")
def test_missing_official_data_keeps_previous_day_with_preparing_notice(self):
today = date.today()
previous = (today - timedelta(days=1)).strftime("%Y-%m-%d")
latest = {
"meta": {"trade_date": "2026-09-07", "source": "tushare"},
"meta": {"trade_date": previous, "source": "tushare"},
"overview": {"limit_up_count": 20},
}
harness = SyncHarness(
FakeRealtimeTodayClient(),
latest,
clock=lambda: at_clock(10, 5),
)
payload = harness.sync_dashboard(today)
harness = SyncHarness(FakeMissingDailyClient(), latest)
payload = harness.sync_dashboard(today.strftime("%Y%m%d"))
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.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 = {
"meta": {"trade_date": "2026-09-07", "source": "tushare"},
"overview": {"limit_up_count": 20},
}
harness = SyncHarness(
FakeMissingDailyClient(),
latest,
clock=lambda: at_clock(10, 5),
)
with self.assertRaises(ValueError) as ctx:
harness.sync_dashboard(today)
self.assertIn("当天盘中行情", str(ctx.exception))
self.assertFalse(harness.database.saved)
def test_intraday_keeps_existing_today_snapshot_when_refresh_fails(self):
today = TRADE_DAY.strftime("%Y%m%d")
existing = {
"meta": {
"trade_date": "2026-09-08",
"realtime": True,
"mode": "realtime",
"source": "tushare",
},
"overview": {"limit_up_count": 11},
"limits": [{"code": "600000"}],
"broken": [],
"down_limits": [],
"yesterday_limits": [],
}
harness = SyncHarness(
FakeMissingDailyClient(),
clock=lambda: at_clock(10, 5),
)
harness.database.get_snapshot = lambda *_args, **_kwargs: copy.deepcopy(existing)
payload = harness.sync_dashboard(today)
meta = payload["meta"]
self.assertEqual(str(meta["trade_date"]).replace("-", ""), today)
self.assertTrue(meta["realtime"])
self.assertEqual(meta["data_status"], "intraday")
self.assertFalse(meta.get("carried_forward"))
def test_lunch_and_after_hours_keep_today_until_official_arrives(self):
today = TRADE_DAY.strftime("%Y%m%d")
for clock in (lambda: at_clock(12, 0), lambda: at_clock(16, 10)):
harness = SyncHarness(
FakeRealtimeTodayClient(),
clock=clock,
)
payload = harness.sync_dashboard(today)
self.assertEqual(str(payload["meta"]["trade_date"]).replace("-", ""), today)
self.assertFalse(payload["meta"].get("carried_forward"))
def test_preopen_and_weekend_still_carry_last_session(self):
latest = {
"meta": {"trade_date": "2026-09-07", "source": "tushare"},
"overview": {"limit_up_count": 20},
}
preopen = SyncHarness(
FakeMissingDailyClient(),
latest,
clock=lambda: at_clock(8, 30),
)
preopen_payload = preopen.sync_dashboard(TRADE_DAY.strftime("%Y%m%d"))
self.assertTrue(preopen_payload["meta"]["carried_forward"])
self.assertEqual(preopen_payload["meta"]["data_status"], "preparing")
self.assertIn("今日数据正在准备,当前展示", preopen_payload["meta"]["display_notice"])
weekend = SyncHarness(
FakeMissingDailyClient(open_today=False),
latest,
clock=lambda: at_clock(10, 5, date(2026, 9, 5)),
)
weekend_payload = weekend.sync_dashboard("20260905")
self.assertTrue(weekend_payload["meta"]["carried_forward"])
def test_history_date_still_uses_official_or_preparing_notice(self):
latest = {
"meta": {"trade_date": "2026-09-01", "source": "tushare"},
"overview": {"limit_up_count": 8},
}
harness = SyncHarness(
FakeMissingDailyClient(),
latest,
clock=lambda: at_clock(10, 5),
)
payload = harness.sync_dashboard("20260902")
self.assertTrue(payload["meta"]["carried_forward"])
self.assertIn("所选日期数据尚未到齐", payload["meta"]["display_notice"])
def test_carried_today_snapshot_is_retried_immediately_in_session(self):
today = TRADE_DAY.strftime("%Y%m%d")
snapshot = {
"meta": {
"source": "tushare",
"trade_date": "2026-09-07",
"carried_forward": True,
"requested_date": "2026-09-08",
"updated_at": at_clock(10, 0).isoformat(),
},
"overview": {"limit_up_count": 1},
}
harness = SyncHarness(
FakeRealtimeTodayClient(),
clock=lambda: at_clock(10, 5),
)
harness.database.get_snapshot = lambda *_args, **_kwargs: copy.deepcopy(snapshot)
payload = harness.get_dashboard(today)
self.assertEqual(str(payload["meta"]["trade_date"]).replace("-", ""), today)
self.assertEqual(payload["meta"]["data_status"], "intraday")
self.assertTrue(harness.database.saved)
self.assertTrue(meta["carried_forward"])
self.assertEqual(meta["data_status"], "preparing")
self.assertIn("今日数据正在准备,当前展示", meta["display_notice"])
self.assertIn("", meta["display_notice"])
self.assertNotIn("No daily data", meta["display_notice"])
self.assertNotEqual(verified_dashboard_result(payload).get("status"), "failed")
def test_weekend_carry_is_not_labeled_as_preparing(self):
snapshot = {
@@ -405,43 +200,19 @@ class DashboardFreshnessTests(unittest.TestCase):
{"meta": {"trade_date": iso, "limit_data_source": "derived"}},
)
now = datetime.now().astimezone().time().replace(tzinfo=None)
if dt_time(15, 5) <= now < dt_time(22, 0):
if datetime.strptime("15:05", "%H:%M").time() <= now < datetime.strptime("22:00", "%H:%M").time():
self.assertFalse(due)
self.assertTrue(derived_due)
else:
self.assertFalse(due)
self.assertFalse(derived_due)
def test_official_catchup_is_due_for_intraday_snapshot_after_close(self):
today = TRADE_DAY.strftime("%Y%m%d")
snapshot = {
"meta": {
"trade_date": "2026-09-08",
"realtime": True,
"mode": "realtime",
}
}
with patch("backend.jobs.refresh.datetime") as mocked:
mocked.now.return_value = at_clock(16, 10)
mocked.strptime = datetime.strptime
self.assertTrue(official_catchup_due(today, snapshot))
official = {
"meta": {
"trade_date": "2026-09-08",
"limit_data_source": "official",
"realtime": False,
}
}
self.assertFalse(official_catchup_due(today, official))
class FrontendRefreshCopyTests(unittest.TestCase):
def test_dashboard_script_distinguishes_partial_from_failure(self):
script = (Path(__file__).resolve().parents[1] / "frontend" / "shared" / "dashboard.js").read_text(encoding="utf-8")
self.assertIn("今日数据正在准备,当前展示", script)
self.assertIn("部分正式数据尚未到齐", script)
self.assertIn("盘中行情", script)
self.assertIn("meta.realtime && actualCompact === requestedCompact", script)
self.assertIn('job.status === "failed"', script)
failed_block = script.split("if (job.status === \"failed\")", 1)[1].split("const query", 1)[0]
self.assertIn("后台刷新失败", failed_block)
+1 -167
View File
@@ -2,8 +2,7 @@ from __future__ import annotations
import unittest
from backend.data.providers.ifind_client import IfindHttpClient
from backend.features.market.charts import ChartDataError, EastmoneyChartClient, HIS_TRENDS_URL, MarketChartClient, TRENDS_URL
from backend.features.market.charts import ChartDataError, EastmoneyChartClient
from server import DashboardService
@@ -73,171 +72,6 @@ class ChartDataProviderTests(unittest.TestCase):
self.client.stock_intraday("abc")
class LookbackChartClient(EastmoneyChartClient):
def __init__(self) -> None:
super().__init__(cache_ttl_seconds=20)
self.requests: list[tuple[str, dict[str, str]]] = []
def _request_json(self, url, params, referer):
self.requests.append((url, params))
if url == TRENDS_URL and params.get("ndays") == "1":
return {"data": {"code": "601318", "name": "中国平安", "preClose": 56.0, "trends": []}}
if url == TRENDS_URL and params.get("ndays") == "5":
return {"data": {"code": "601318", "name": "中国平安", "preClose": 56.0, "trends": []}}
if url == HIS_TRENDS_URL:
return {
"data": {
"code": "601318",
"name": "中国平安",
"preClose": 55.8,
"trends": [
"2026-09-07 09:30,55.80,55.90,56.00,55.70,100,5580.00,55.900",
"2026-09-07 15:00,56.10,56.20,56.30,56.00,200,11240.00,56.150",
"2026-09-08 09:30,0,0,0,0,0,0.00,0",
],
}
}
raise ChartDataError("unexpected url")
class ChartLookbackTests(unittest.TestCase):
def setUp(self) -> None:
EastmoneyChartClient._cache.clear()
self.client = LookbackChartClient()
def test_empty_today_falls_back_to_latest_available_session(self):
payload = self.client.stock_intraday("601318")
urls = [url for url, _ in self.client.requests]
self.assertEqual(urls[0], TRENDS_URL)
self.assertEqual(self.client.requests[0][1]["ndays"], "1")
self.assertEqual(urls[1], TRENDS_URL)
self.assertEqual(self.client.requests[1][1]["ndays"], "5")
self.assertEqual(urls[2], HIS_TRENDS_URL)
self.assertEqual(payload["trade_date"], "2026-09-07")
self.assertEqual([point["time"] for point in payload["points"]], ["09:30", "15:00"])
self.assertEqual(payload["points"][0]["close"], 55.9)
def test_delay_multiday_can_recover_without_his(self):
class DelayFive(EastmoneyChartClient):
def __init__(self):
super().__init__(cache_ttl_seconds=20)
self.requests = []
def _request_json(self, url, params, referer):
self.requests.append((url, params))
if params.get("ndays") == "1":
return {"data": {"code": "000001", "name": "平安银行", "preClose": 11.7, "trends": []}}
return {
"data": {
"code": "000001",
"name": "平安银行",
"preClose": 11.5,
"trends": [
"2026-09-07 09:30,11.50,11.60,11.70,11.40,100,1160.00,11.600",
"2026-09-07 15:00,11.70,11.80,11.90,11.60,200,2360.00,11.750",
],
}
}
EastmoneyChartClient._cache.clear()
client = DelayFive()
payload = client.stock_intraday("000001")
self.assertEqual(payload["trade_date"], "2026-09-07")
self.assertEqual(len(payload["points"]), 2)
self.assertEqual([url for url, _ in client.requests], [TRENDS_URL, TRENDS_URL])
def test_sh_sz_cyb_codes_use_correct_secid(self):
for code, secid in (("601318", "1.601318"), ("000001", "0.000001"), ("300750", "0.300750")):
EastmoneyChartClient._cache.clear()
client = LookbackChartClient()
client.stock_intraday(code)
self.assertEqual(client.requests[0][1]["secid"], secid)
class FakeHub:
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)
if self.error:
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:
EastmoneyChartClient._cache.clear()
def test_datahub_success_skips_old_channel(self):
hub = FakeHub(
{
"entity_type": "stock",
"identifier": "601318",
"name": "中国平安",
"code": "601318",
"trade_date": "2026-09-08",
"previous_close": 56.36,
"points": [{"date": "2026-09-08", "time": "09:30", "close": 56.5, "average": 56.4}],
"source": "datahub",
}
)
fallback = LookbackChartClient()
client = MarketChartClient(IfindHttpClient(), fallback, hub)
payload = client.stock_intraday("601318")
self.assertEqual(payload["source"], "datahub")
self.assertEqual(hub.calls, ["601318"])
self.assertEqual(fallback.requests, [])
def test_datahub_timeout_or_empty_falls_back_to_eastmoney(self):
fallback = LookbackChartClient()
for hub in (
FakeHub(chart=None),
FakeHub(error=RuntimeError("timeout")),
FakeHub(error=RuntimeError("datahub exploded")),
FakeHub(chart={"points": []}),
):
EastmoneyChartClient._cache.clear()
fallback.requests.clear()
client = MarketChartClient(IfindHttpClient(), fallback, hub)
payload = client.stock_intraday("000001")
self.assertEqual(payload["trade_date"], "2026-09-07")
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
def _payload(code: str, name: str):
+6 -257
View File
@@ -12,7 +12,6 @@ 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]
@@ -65,11 +64,9 @@ class FakeClient(DatahubClient):
meta={"tier": "official", "trade_date": "20240902", "stale": False, "staleness_seconds": 0},
)
self.paths: list[str] = []
self.calls: list[tuple[str, dict[str, Any]]] = []
def get(self, path: str, params: dict[str, Any] | None = None) -> DatahubResponse:
self.paths.append(path)
self.calls.append((path, {key: value for key, value in (params or {}).items()}))
if TOKEN in json.dumps(params or {}) or TOKEN in path:
raise AssertionError("token leaked into url")
if self.error:
@@ -85,21 +82,17 @@ def flags(**enabled: tuple[bool, bool]) -> DatahubSettings:
class DatahubBridgeTests(unittest.TestCase):
def setUp(self) -> None:
LEDGER.clear()
def test_default_config_enables_official_reads(self) -> None:
def test_default_config_keeps_legacy_and_does_not_call_datahub(self) -> None:
settings = DatahubSettings.load(environ={}, credentials={})
self.assertTrue(settings.any_enabled())
self.assertTrue(all(settings.flags(name).read and not settings.flags(name).shadow for name in DATASETS))
client = FakeClient()
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"))
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, ["/v1/bars/daily"])
self.assertEqual(legacy.calls, [])
self.assertEqual(LEDGER.snapshot()[0]["route"], "datahub")
self.assertEqual(client.paths, [])
self.assertEqual(len(legacy.calls), 1)
def test_each_dataset_has_independent_read_flag(self) -> None:
settings = flags(daily=(True, False), auction=(False, False))
@@ -109,13 +102,6 @@ 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]] = []
@@ -363,243 +349,6 @@ class DatahubBridgeTests(unittest.TestCase):
self.assertEqual(rows[0]["amount"], 2000.0)
self.assertEqual(len(legacy.calls), 1)
def test_try_intraday_respects_switch_and_falls_back_on_bad_payload(self) -> None:
closed = DatahubBridge(flags(), FakeClient(error=DatahubError("INTERNAL", "should not run")))
self.assertIsNone(closed.try_intraday("601318"))
empty = DatahubBridge(
flags(intraday=(True, False)),
FakeClient(response=DatahubResponse(data={"points": []}, meta={"stale": False})),
)
self.assertIsNone(empty.try_intraday("601318"))
stale = DatahubBridge(
flags(intraday=(True, False)),
FakeClient(response=DatahubResponse(
data={
"entity_type": "stock",
"code": "601318",
"trade_date": "2026-09-07",
"previous_close": 55.8,
"points": [{"date": "2026-09-07", "time": "09:30", "close": 55.9, "avg_price": 55.85}],
},
meta={"stale": True},
)),
)
self.assertIsNone(stale.try_intraday("601318"))
ok = DatahubBridge(
flags(intraday=(True, False)),
FakeClient(response=DatahubResponse(
data={
"entity_type": "stock",
"identifier": "601318",
"name": "中国平安",
"code": "601318",
"trade_date": "2026-09-08",
"previous_close": 56.36,
"points": [
{"date": "2026-09-08", "time": "09:30", "close": 0},
{"date": "2026-09-08", "time": "09:31", "close": 56.5, "avg_price": 56.4},
],
},
meta={"stale": False},
)),
)
chart = ok.try_intraday("601318")
self.assertEqual(chart["source"], "datahub")
self.assertEqual(len(chart["points"]), 1)
self.assertEqual(chart["points"][0]["average"], 56.4)
self.assertEqual(ok.client.paths, ["/v1/intraday/points"])
self.assertEqual(ok.client.calls, [("/v1/intraday/points", {"code": "601318"})])
self.assertNotIn("date", ok.client.calls[0][1])
timeout = DatahubBridge(
flags(intraday=(True, False)),
FakeClient(error=DatahubError("TIMEOUT", "datahub request timed out")),
)
self.assertIsNone(timeout.try_intraday("601318"))
broken = DatahubBridge(
flags(intraday=(True, False)),
FakeClient(error=DatahubError("INTERNAL", "datahub exploded")),
)
self.assertIsNone(broken.try_intraday("601318"))
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_try_daily_chart_keeps_usable_bars_when_coverage_incomplete(self) -> None:
rows = [
{
"ts_code": "000001.SZ",
"trade_date": "20240901",
"open": 10.0,
"high": 10.4,
"low": 9.9,
"close": 10.2,
"volume": 100000,
"amount": 2000000,
},
{
"ts_code": "000001.SZ",
"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,
"incomplete": True,
"coverage": {"complete": False, "missing_count": 127},
"source": "tushare:daily",
},
)
),
)
chart = hub.try_daily_chart("000001.SZ", "20240902", 90, "daily")
self.assertIsNotNone(chart)
self.assertEqual(chart[-1]["trade_date"], "2024-09-02")
self.assertEqual(chart[-1]["close"], 10.4)
def test_gateway_tushare_assembly_binds_hooks_on_inner_client(self) -> None:
quotes = [
{
"ts_code": f"{index:06d}.SZ",
"name": f"S{index}",
"pre_close": 10.0,
"open": 10.0,
"high": 10.5,
"low": 9.8,
"close": 10.2,
"vol": 100.0,
"amount": 1000.0,
"quote_date": "20240902",
}
for index in range(1, 221)
]
hub_client = FakeClient(
response=DatahubResponse(
data=quotes,
meta={"stale": False, "staleness_seconds": 0, "source": "eastmoney_clist"},
)
)
gateway = build_data_gateway(
{"tushare_token": "tok"},
datahub_settings=flags(quotes=(True, False), daily=(True, False)),
)
gateway.datahub.client = hub_client
wrapped = gateway.tushare()
inner = wrapped._legacy
self.assertTrue(callable(getattr(inner, "try_market_quotes", None)))
self.assertTrue(callable(getattr(inner, "try_index_quotes", None)))
self.assertTrue(callable(getattr(inner, "record_datahub_legacy", None)))
self.assertIs(inner.query.__self__, wrapped)
self.assertEqual(inner.query.__func__, wrapped.query.__func__)
self.assertFalse(hasattr(type(inner), "try_market_quotes"))
rows = inner.try_market_quotes("20240902")
self.assertGreaterEqual(len(rows or []), 200)
self.assertIn("/v1/quotes/latest", hub_client.paths)
hub_client.response = DatahubResponse(
data=[dict(HUB_DAILY)],
meta={"stale": False, "staleness_seconds": 0, "source": "tushare:daily"},
)
daily = inner.query("daily", {"trade_date": "20240902"}, "ts_code,amount")
self.assertEqual(daily[0]["amount"], 2000.0)
self.assertIn("/v1/bars/daily", hub_client.paths)
def test_features_do_not_import_datahub_client(self) -> None:
violations = []
for path in (ROOT / "backend" / "features").rglob("*.py"):
+1 -83
View File
@@ -3,10 +3,9 @@ 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 RealtimeAggregateError, WebRealtimeAggregator
from backend.data.realtime import WebRealtimeAggregator
from backend.features.heaven.engine import _market_line_scores, build_manual_market_hexagram
from server import DashboardService
from backend.data.providers.tushare_client import (
@@ -378,87 +377,6 @@ 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")
@patch.object(WebRealtimeAggregator, "_get_text")
def test_tencent_stock_quote_keeps_expected_date(self, get_text: MagicMock):
fields = [""] * 38
fields[1] = "浦发银行"
fields[2] = "600000"
fields[3] = "11.20"
fields[4] = "11.00"
fields[5] = "11.10"
fields[6] = "1234"
fields[30] = "20260720103000"
fields[33] = "11.30"
fields[34] = "11.00"
fields[37] = "1380"
get_text.return_value = (f'v_sh600000="{"~".join(fields)}";', 0)
quote = WebRealtimeAggregator().tencent_stock_quote("600000", "20260720")
self.assertEqual(quote["ts_code"], "600000.SH")
self.assertEqual(quote["quote_date"], "20260720")
self.assertEqual(quote["vol"], 123400)
self.assertAlmostEqual(quote["amount"], 13_800_000)
@patch.object(WebRealtimeAggregator, "_get_json")
def test_eastmoney_stock_quote_rejects_stale_date(self, get_json: MagicMock):
epoch = datetime(2026, 7, 19, 15, 0).timestamp()
get_json.return_value = {
"rc": 0,
"data": {
"f43": 11.2,
"f44": 11.3,
"f45": 11.0,
"f46": 11.1,
"f47": 10,
"f48": 50000000,
"f57": "300750",
"f58": "宁德时代",
"f60": 11.0,
"f86": epoch,
},
}
with self.assertRaises(RealtimeAggregateError):
WebRealtimeAggregator().eastmoney_stock_quote("300750.SZ", "20260720")
if __name__ == "__main__":
unittest.main()
+1 -1
View File
@@ -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), 28)
self.assertEqual(len(claimed), 27)
if __name__ == "__main__":
-333
View File
@@ -1,16 +1,8 @@
from __future__ import annotations
import unittest
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):
@@ -89,65 +81,6 @@ 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()
@@ -197,272 +130,6 @@ class RealtimeDashboardTests(unittest.TestCase):
self.assertEqual(dashboard["meta"]["limit_data_source"], "derived")
self.assertIn("日线数据推算", dashboard["meta"]["notice"])
def test_calendar_open_flag_accepts_string_and_bool(self):
self.assertTrue(calendar_is_open(1))
self.assertTrue(calendar_is_open("1"))
self.assertTrue(calendar_is_open(True))
self.assertFalse(calendar_is_open(0))
self.assertFalse(calendar_is_open("0"))
self.assertFalse(calendar_is_open(False))
original_query = self.client.query
def query(api_name, params=None, fields=""):
if api_name == "trade_cal":
return [
{
"cal_date": params.get("start_date"),
"is_open": "1",
"pretrade_date": "20260907",
}
]
return original_query(api_name, params, fields)
self.client.query = query
trade_date, previous = self.client.resolve_trade_context("20260908")
self.assertEqual(trade_date, "20260908")
self.assertEqual(previous, "20260907")
def test_session_clock_uses_realtime_until_official_window(self):
today = "20260908"
self.client.clock = lambda: datetime(
2026, 9, 8, 10, 5, tzinfo=timezone(timedelta(hours=8))
)
self.assertTrue(self.client.should_use_realtime(today, today))
self.client.clock = lambda: datetime(
2026, 9, 8, 16, 10, tzinfo=timezone(timedelta(hours=8))
)
self.assertFalse(self.client.should_use_realtime(today, today))
def test_realtime_dashboard_survives_missing_limit_table(self):
original_query = self.client.query
def query(api_name, params=None, fields=""):
if api_name == "stk_limit":
return []
return original_query(api_name, params, fields)
self.client.query = query
TushareClient._realtime_reference_cache.clear()
dashboard = self.client._realtime_dashboard("20260720", "20260720", "20260717")
self.assertTrue(dashboard["meta"]["realtime"])
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")
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"])
def test_gateway_dashboard_uses_bound_market_quotes(self) -> None:
from backend.data import build_data_gateway
from backend.data.datahub.client import DatahubResponse
from backend.data.datahub.settings import DATASETS, DatahubSettings, DatasetFlags
from backend.data.gateway import DataGateway
from backend.data.providers.tushare import TushareProvider
quotes = [
{
"ts_code": item["ts_code"],
"name": item["name"],
"pre_close": item["pre_close"],
"open": item["open"],
"high": item["high"],
"low": item["low"],
"close": item["close"],
"vol": item["vol"],
"amount": item["amount"],
"quote_date": "20260720",
}
for item in FREE_QUOTES
]
extras = [
{
"ts_code": f"{index:06d}.SZ",
"name": f"X{index}",
"pre_close": 10.0,
"open": 10.0,
"high": 10.2,
"low": 9.8,
"close": 10.1,
"vol": 100.0,
"amount": 1000.0,
"quote_date": "20260720",
}
for index in range(10, 230)
]
class QuoteHub:
def __init__(self):
self.calls = []
def quotes_latest(self, **params):
return self.get("/v1/quotes/latest", params)
def get(self, path, params=None):
self.calls.append(path)
if path == "/v1/quotes/latest":
return DatahubResponse(
data=quotes + extras,
meta={"stale": False, "staleness_seconds": 0, "source": "eastmoney_clist"},
)
raise AssertionError(path)
datasets = {name: DatasetFlags(name) for name in DATASETS}
datasets["quotes"] = DatasetFlags("quotes", read=True, shadow=False)
settings = DatahubSettings(base_url="http://127.0.0.1:9", token="tok", datasets=datasets)
base = build_data_gateway({"tushare_token": "tok"}, datahub_settings=settings)
gateway = DataGateway(
policy=base.policy,
quality=base.quality,
tushare_provider=TushareProvider(
lambda: "tok",
client_factory=lambda token: FakeRealtimeClient(token),
),
ifind_provider=base.ifind_provider,
chart_data=base.chart_data,
realtime_observer=base.realtime_observer,
datahub=base.datahub,
)
gateway.datahub.client = QuoteHub()
wrapped = gateway.tushare()
inner = wrapped._legacy
inner.clock = lambda: datetime(2026, 7, 20, 10, 30, tzinfo=timezone(timedelta(hours=8)))
inner.realtime_aggregator = FakeFreeAggregator(fail=True)
TushareClient._realtime_reference_cache.clear()
dashboard = wrapped.dashboard("20260720")
self.assertEqual(dashboard["meta"]["quote_source"], "datahub")
self.assertIn("/v1/quotes/latest", gateway.datahub.client.calls)
self.assertTrue(callable(getattr(inner, "try_market_quotes", None)))
self.assertFalse(hasattr(type(inner), "try_market_quotes"))
if __name__ == "__main__":
unittest.main()
-305
View File
@@ -5,10 +5,6 @@ import unittest
from datetime import datetime, timedelta
from unittest.mock import patch
from backend.data.providers.tushare_client import TushareError
from backend.data.realtime import RealtimeAggregateError
from backend.features.market.charts import ChartDataError
from backend.features.market.service import TODAY_DAILY_UNAVAILABLE_NOTICE
from server import DashboardService
@@ -21,10 +17,6 @@ class DetailDatabaseStub:
def list_notes(user_id, code=""):
return []
@staticmethod
def get_snapshot(trade_date):
return {}
class RealtimeClientStub:
quote_calls = 0
@@ -69,122 +61,6 @@ class FixedPreopenDatetime(datetime):
return cls.fixed_now
class FixedLunchDatetime(datetime):
fixed_now = datetime(2026, 7, 31, 11, 45).astimezone()
@classmethod
def now(cls, tz=None):
return cls.fixed_now
class FixedAfterCloseDatetime(datetime):
fixed_now = datetime(2026, 7, 31, 15, 30).astimezone()
@classmethod
def now(cls, tz=None):
return cls.fixed_now
class DeniedRealtimeClientStub:
quote_calls = 0
def __init__(self, token):
self.token = token
@staticmethod
def resolve_trade_context(requested_date):
return requested_date, requested_date
@classmethod
def realtime_stock_quote(cls, ts_code, reference_date=""):
cls.quote_calls += 1
raise TushareError("没有接口访问权限")
class FreeQuoteAggregator:
def __init__(self, quote=None, fail=False):
self.quote = quote
self.fail = fail
self.tencent_calls = 0
self.eastmoney_calls = 0
def tencent_stock_quote(self, code, expected_date=""):
self.tencent_calls += 1
if self.fail:
raise RealtimeAggregateError("tencent down")
if self.quote and self.quote.get("source") == "eastmoney_stock":
raise RealtimeAggregateError("tencent empty")
if self.quote:
return self.quote
raise RealtimeAggregateError("tencent empty")
def eastmoney_stock_quote(self, code, expected_date=""):
self.eastmoney_calls += 1
if self.fail:
raise RealtimeAggregateError("eastmoney down")
if self.quote and self.quote.get("source") == "eastmoney_stock":
return self.quote
raise RealtimeAggregateError("eastmoney empty")
class IntradayChartStub:
def __init__(self, points, previous_close=10.0, trade_date="2026-07-31"):
self.points = points
self.previous_close = previous_close
self.trade_date = trade_date
def stock_daily(self, code, end_date, limit=90):
raise ChartDataError("iFinD daily unavailable")
def stock_intraday(self, code):
return {
"trade_date": self.trade_date,
"previous_close": self.previous_close,
"points": self.points,
}
def _history_payload(code="002141"):
yesterday = (FixedMarketDatetime.fixed_now - timedelta(days=1)).strftime("%Y-%m-%d")
return {
"meta": {"trade_date": yesterday, "source": "tushare"},
"stock": {"code": code, "name": "旧名称", "price": 10, "change": 7.1},
"prices": [
{
"trade_date": yesterday,
"open": 9.5,
"high": 10.1,
"low": 9.4,
"close": 10,
"change": 7.1,
"volume": 100,
"amount_billion": 1.1,
}
],
"moneyflow": {},
}
def _free_quote(source="tencent_qt", **overrides):
quote = {
"ts_code": "002141.SZ",
"name": "贤程科技",
"pre_close": 10.0,
"open": 10.2,
"high": 10.8,
"low": 10.1,
"close": 10.6,
"vol": 250000,
"amount": 26_500_000,
"quote_date": "20260731",
"quote_time_epoch": int(datetime(2026, 7, 31, 10, 31).timestamp()),
"source": source,
"turnover_rate": 2.5,
}
quote.update(overrides)
return quote
class StockDetailRealtimeTests(unittest.TestCase):
def setUp(self):
self.service = DashboardService.__new__(DashboardService)
@@ -192,11 +68,7 @@ class StockDetailRealtimeTests(unittest.TestCase):
self.service.database = DetailDatabaseStub()
self.service._request_context = threading.local()
self.service._request_context.user_id = 1
self.service.ifind = None
self.service.realtime_aggregator = None
self.service.chart_data = None
RealtimeClientStub.quote_calls = 0
DeniedRealtimeClientStub.quote_calls = 0
def test_today_detail_merges_rt_quote_without_mutating_daily_cache(self):
today = FixedMarketDatetime.fixed_now.strftime("%Y%m%d")
@@ -290,183 +162,6 @@ class StockDetailRealtimeTests(unittest.TestCase):
self.assertEqual(result["stock"]["change"], 1.2)
self.assertEqual(RealtimeClientStub.quote_calls, 0)
def test_today_detail_falls_back_to_tencent_quote_when_rt_k_denied(self):
today = FixedMarketDatetime.fixed_now.strftime("%Y%m%d")
aggregator = FreeQuoteAggregator(_free_quote())
self.service.realtime_aggregator = aggregator
DeniedRealtimeClientStub.quote_calls = 0
with patch("backend.features.market.service.datetime", FixedMarketDatetime), patch(
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
):
result = self.service._prepare_stock_detail(_history_payload(), "002141", today)
bar = result["prices"][-1]
self.assertEqual(bar["trade_date"], "2026-07-31")
self.assertTrue(bar["realtime"])
self.assertEqual(bar["open"], 10.2)
self.assertEqual(bar["high"], 10.8)
self.assertEqual(bar["low"], 10.1)
self.assertEqual(bar["close"], 10.6)
self.assertAlmostEqual(bar["change"], 6.0, places=4)
self.assertEqual(bar["volume"], 2500)
self.assertAlmostEqual(bar["amount_billion"], 0.265)
self.assertEqual(len(result["prices"]), 2)
self.assertEqual(result["meta"]["notice"], "")
self.assertEqual(aggregator.tencent_calls, 1)
self.assertEqual(DeniedRealtimeClientStub.quote_calls, 1)
def test_today_detail_falls_back_to_eastmoney_then_intraday(self):
today = FixedMarketDatetime.fixed_now.strftime("%Y%m%d")
aggregator = FreeQuoteAggregator(
_free_quote("eastmoney_stock", ts_code="600000.SH", name="浦发银行"),
)
self.service.realtime_aggregator = aggregator
DeniedRealtimeClientStub.quote_calls = 0
with patch("backend.features.market.service.datetime", FixedMarketDatetime), patch(
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
):
result = self.service._prepare_stock_detail(_history_payload("600000"), "600000", today)
self.assertEqual(result["prices"][-1]["trade_date"], "2026-07-31")
self.assertEqual(result["prices"][-1]["close"], 10.6)
self.assertEqual(aggregator.tencent_calls, 1)
self.assertEqual(aggregator.eastmoney_calls, 1)
aggregator = FreeQuoteAggregator(fail=True)
self.service.realtime_aggregator = aggregator
self.service.chart_data = IntradayChartStub(
[
{
"date": "2026-07-31",
"time": "09:30",
"open": 10.1,
"high": 10.2,
"low": 10.0,
"close": 10.15,
"volume": 120,
"amount": 121800,
},
{
"date": "2026-07-31",
"time": "10:05",
"open": 10.15,
"high": 10.5,
"low": 9.9,
"close": 10.4,
"volume": 80,
"amount": 83200,
},
]
)
with patch("backend.features.market.service.datetime", FixedMarketDatetime), patch(
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
):
result = self.service._prepare_stock_detail(_history_payload("300750"), "300750", today)
bar = result["prices"][-1]
self.assertEqual(bar["trade_date"], "2026-07-31")
self.assertEqual(bar["open"], 10.1)
self.assertEqual(bar["high"], 10.5)
self.assertEqual(bar["low"], 9.9)
self.assertEqual(bar["close"], 10.4)
self.assertAlmostEqual(bar["change"], 4.0, places=4)
self.assertEqual(bar["volume"], 200)
self.assertTrue(bar["realtime"])
def test_today_detail_keeps_history_when_free_sources_fail(self):
today = FixedMarketDatetime.fixed_now.strftime("%Y%m%d")
self.service.realtime_aggregator = FreeQuoteAggregator(fail=True)
self.service.chart_data = IntradayChartStub([], trade_date="2026-07-30")
DeniedRealtimeClientStub.quote_calls = 0
with patch("backend.features.market.service.datetime", FixedMarketDatetime), patch(
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
):
result = self.service._prepare_stock_detail(_history_payload(), "002141", today)
self.assertEqual(result["prices"][-1]["trade_date"], "2026-07-30")
self.assertFalse(result["meta"].get("realtime", False))
self.assertEqual(result["meta"]["notice"], TODAY_DAILY_UNAVAILABLE_NOTICE)
self.assertEqual(len(result["prices"]), 1)
def test_lunch_keeps_morning_realtime_bar(self):
today = FixedLunchDatetime.fixed_now.strftime("%Y%m%d")
self.service.realtime_aggregator = FreeQuoteAggregator(
_free_quote(quote_time_epoch=int(datetime(2026, 7, 31, 11, 30).timestamp()))
)
DeniedRealtimeClientStub.quote_calls = 0
with patch("backend.features.market.service.datetime", FixedLunchDatetime), patch(
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
):
result = self.service._prepare_stock_detail(_history_payload(), "002141", today)
self.assertEqual(result["prices"][-1]["trade_date"], "2026-07-31")
self.assertTrue(result["meta"]["realtime"])
def test_after_close_keeps_forming_bar_until_official_ready(self):
today = FixedAfterCloseDatetime.fixed_now.strftime("%Y%m%d")
self.service.realtime_aggregator = FreeQuoteAggregator(_free_quote())
DeniedRealtimeClientStub.quote_calls = 0
with patch("backend.features.market.service.datetime", FixedAfterCloseDatetime), patch(
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
):
forming = self.service._prepare_stock_detail(_history_payload(), "002141", today)
self.assertEqual(forming["prices"][-1]["trade_date"], "2026-07-31")
self.assertTrue(forming["prices"][-1]["realtime"])
official = _history_payload()
official["prices"].append(
{
"trade_date": "2026-07-31",
"open": 10.15,
"high": 10.9,
"low": 10.05,
"close": 10.7,
"change": 7.0,
"volume": 1800,
"amount_billion": 0.3,
}
)
RealtimeClientStub.quote_calls = 0
with patch("backend.features.market.service.datetime", FixedAfterCloseDatetime), patch(
"backend.features.market.service.TushareClient", RealtimeClientStub
):
replaced = self.service._prepare_stock_detail(official, "002141", today)
self.assertEqual(replaced["prices"][-1]["close"], 10.7)
self.assertFalse(replaced["prices"][-1].get("realtime", False))
self.assertEqual(len(replaced["prices"]), 2)
self.assertEqual(RealtimeClientStub.quote_calls, 0)
def test_same_date_bar_is_replaced_not_duplicated(self):
today = FixedMarketDatetime.fixed_now.strftime("%Y%m%d")
payload = _history_payload()
payload["prices"].append(
{
"trade_date": "2026-07-31",
"open": 10.0,
"high": 10.1,
"low": 9.9,
"close": 10.05,
"change": 0.5,
"volume": 10,
"amount_billion": 0.01,
"realtime": True,
}
)
self.service.realtime_aggregator = FreeQuoteAggregator(_free_quote())
DeniedRealtimeClientStub.quote_calls = 0
with patch("backend.features.market.service.datetime", FixedMarketDatetime), patch(
"backend.features.market.service.TushareClient", DeniedRealtimeClientStub
):
result = self.service._prepare_stock_detail(payload, "002141", today)
self.assertEqual(len(result["prices"]), 2)
self.assertEqual(result["prices"][-1]["close"], 10.6)
self.assertEqual(result["prices"][-1]["trade_date"], "2026-07-31")
if __name__ == "__main__":
unittest.main()
+2 -2
View File
@@ -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 and intraday dashboard fallback"},
{"provider": "tencent", "path": "backend/data/realtime.py", "runtime_role": "index observation and intraday quote 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_domains": [
{"provider": "tushare", "path": "backend/data/providers/tushare_transport.py", "responsibility": "HTTP transport and provider errors"},
+1 -1
View File
@@ -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` 不传 codes 即全市场,`/v1/indexes/quotes` `/v1/intraday/points`);永不写入 eod_* 正式表
- 盘中观察(provisional):东财/腾讯指数报价、个股最新价、分时点(`/v1/quotes/latest` `/v1/indexes/quotes` `/v1/intraday/points`);永不写入 eod_* 正式表
- 暂存 → 校验 → 整批原子发布 → 可回滚
- `/v1` 稳定接口(`X-Datahub-Token`
- `/admin/` 最小管理后台(总览 / 数据源 / 调度 / 发布 / 数据集 / 审计)
+21 -153
View File
@@ -13,17 +13,7 @@ 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 = (
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
"AppleWebKit/537.36 (KHTML, like Gecko) Chrome/138.0.0.0 Safari/537.36"
@@ -68,11 +58,7 @@ class EastmoneyAdapter(MarketAdapter):
codes = params.get("codes") or []
if isinstance(codes, str):
codes = [item.strip() for item in codes.split(",") if item.strip()]
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()
return self.fetch_quotes(list(codes))
raise AdapterError(f"{self.name} unsupported dataset: {dataset}")
def normalize(self, dataset: str, rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
@@ -180,59 +166,7 @@ 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]:
def fetch_intraday(self, ts_code: str) -> dict[str, Any]:
code = str(ts_code or "").upper()
if code in INDEX_SECIDS:
secid = INDEX_SECIDS[code]
@@ -244,32 +178,25 @@ class EastmoneyAdapter(MarketAdapter):
secid = f"{market}.{symbol}"
entity = "stock"
identifier = symbol
params = {
"secid": secid,
"fields1": "f1,f2,f3,f4,f5,f6,f7,f8,f9,f10,f11,f12,f13",
"fields2": "f51,f52,f53,f54,f55,f56,f57,f58",
"iscr": "0",
}
data: dict[str, Any] = {}
points: list[dict[str, Any]] = []
last_error: Exception | None = None
for url, ndays in ((TRENDS_URL, "1"), (TRENDS_URL, "5"), (HIS_TRENDS_URL, "5")):
try:
payload = self._get_json(
url,
{**params, "ndays": ndays},
referer="https://quote.eastmoney.com/",
)
except AdapterError as exc:
last_error = exc
continue
data = payload.get("data") or {}
parsed = [point for raw in data.get("trends") or [] if (point := _parse_trend(raw))]
points = _preferred_session(parsed, date)
if points:
break
payload = self._get_json(
TRENDS_URL,
{
"secid": secid,
"fields1": "f1,f2,f3,f4,f5,f6,f7,f8,f9,f10,f11,f12,f13",
"fields2": "f51,f52,f53,f54,f55,f56,f57,f58",
"iscr": "0",
"ndays": "1",
},
referer="https://quote.eastmoney.com/",
)
data = payload.get("data") or {}
points = []
for raw in data.get("trends") or []:
point = _parse_trend(raw)
if point:
points.append(point)
if not points:
raise AdapterError("No intraday chart data returned") from last_error
raise AdapterError("No intraday chart data returned")
return {
"entity_type": entity,
"identifier": identifier,
@@ -300,62 +227,6 @@ class EastmoneyAdapter(MarketAdapter):
raise AdapterError(f"eastmoney request failed: {exc}") from exc
def _preferred_session(points: list[dict[str, Any]], preferred_date: str = "") -> list[dict[str, Any]]:
if not points:
return []
want = ""
digits = str(preferred_date or "").replace("-", "")[:8]
if len(digits) == 8 and digits.isdigit():
want = f"{digits[:4]}-{digits[4:6]}-{digits[6:8]}"
if want:
matched = [point for point in points if str(point.get("date") or "") == want]
if matched:
return matched
latest = max(str(point.get("date") or "") for point in points)
if not latest:
return points
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(",")
@@ -366,14 +237,11 @@ def _parse_trend(raw: Any) -> dict[str, Any] | None:
when = datetime.strptime(stamp, "%Y-%m-%d %H:%M")
except ValueError:
return None
close = round4(finite_number(parts[2]))
if close <= 0:
return None
return {
"time": when.strftime("%H:%M"),
"date": when.strftime("%Y-%m-%d"),
"open": round4(finite_number(parts[1])),
"close": close,
"close": round4(finite_number(parts[2])),
"high": round4(finite_number(parts[3])),
"low": round4(finite_number(parts[4])),
"avg_price": round4(finite_number(parts[7] if len(parts) > 7 else parts[2])),
+3 -75
View File
@@ -14,7 +14,6 @@ from datahub.adapters.eastmoney import EastmoneyAdapter
from datahub.adapters.tencent import TencentAdapter
from datahub.codes import resolve_code
from datahub.db import HubDB
from datahub.governance.lkg import LastKnownGood
from datahub.timeutil import isoformat, now_shanghai, yyyymmdd
QUOTE_TTL = 60
@@ -64,36 +63,9 @@ 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:
return fetch_market_quotes(db)
raise RealtimeApiError("INVALID_ARGUMENT", "codes is required")
resolved: list[str] = []
for code in codes[:60]:
item = resolve_code(db, code) or _guess_ts_code(code)
@@ -136,13 +108,10 @@ def fetch_intraday(db: HubDB, code: str, date: str = "") -> dict[str, Any]:
return cached
adapter = EastmoneyAdapter()
try:
payload_data = adapter.fetch_intraday(ts_code, date)
payload_data = adapter.fetch_intraday(ts_code)
source = "eastmoney:trends2"
except Exception as exc:
recovered = _load_intraday_lkg(db, ts_code, date)
if recovered is None:
raise RealtimeApiError("SOURCE_UNAVAILABLE", f"intraday unavailable: {exc}") from exc
return recovered
raise RealtimeApiError("SOURCE_UNAVAILABLE", f"intraday unavailable: {exc}") from exc
payload = _envelope(
payload_data,
{
@@ -158,47 +127,6 @@ def fetch_intraday(db: HubDB, code: str, date: str = "") -> dict[str, Any]:
return payload
def _load_intraday_lkg(db: HubDB, ts_code: str, date: str = "") -> dict[str, Any] | None:
store = LastKnownGood(db)
keys = [f"intraday:{ts_code}:{date or 'today'}"]
if date:
keys.append(f"intraday:{ts_code}:today")
for key in keys:
item = store.load(key)
payload = _lkg_payload(item)
if payload is not None:
return payload
row = db.fetchone(
"SELECT * FROM last_known_good WHERE cache_key LIKE ? ORDER BY stored_at DESC LIMIT 1",
(f"intraday:{ts_code}:%",),
)
if not row:
return None
try:
raw = json.loads(row["payload"])
except json.JSONDecodeError:
return None
return _mark_stale(raw) if isinstance(raw, dict) else None
def _lkg_payload(item: dict[str, Any] | None) -> dict[str, Any] | None:
if not item:
return None
payload = item.get("payload")
return _mark_stale(payload) if isinstance(payload, dict) else None
def _mark_stale(payload: dict[str, Any]) -> dict[str, Any] | None:
data = payload.get("data")
if not isinstance(data, dict) or not data.get("points"):
return None
stamped = dict(payload)
meta = dict(stamped.get("meta") or {})
meta["stale"] = True
stamped["meta"] = meta
return stamped
def _guess_ts_code(code: str) -> str | None:
raw = str(code or "").strip().upper()
if "." in raw:
+3 -10
View File
@@ -247,13 +247,11 @@ class V1API:
)
def quotes_latest(self, q: dict[str, str]) -> dict[str, Any]:
from datahub.realtime_serve import RealtimeApiError, fetch_market_quotes, fetch_quotes
from datahub.realtime_serve import RealtimeApiError, fetch_quotes
codes = [item.strip() for item in str(q.get("codes") or "").split(",") if item.strip()]
try:
if codes:
return fetch_quotes(self.db, codes)
return fetch_market_quotes(self.db)
return fetch_quotes(self.db, codes)
except RealtimeApiError as exc:
raise ApiError(exc.code, exc.message) from exc
@@ -271,13 +269,8 @@ class V1API:
code = str(q.get("code") or "").strip()
if not code:
raise ApiError("INVALID_ARGUMENT", "code is required")
raw_date = str(q.get("date") or "").strip()
try:
trade_date = yyyymmdd(raw_date or now_shanghai())
except ValueError as exc:
raise ApiError("INVALID_ARGUMENT", str(exc)) from exc
try:
return fetch_intraday(self.db, code, trade_date)
return fetch_intraday(self.db, code, yyyymmdd(q.get("date") or ""))
except RealtimeApiError as exc:
raise ApiError(exc.code, exc.message) from exc
@@ -1,232 +0,0 @@
from __future__ import annotations
import tempfile
import unittest
from pathlib import Path
from unittest.mock import patch
from datahub.adapters.base import AdapterError
from datahub.adapters.eastmoney import HIS_TRENDS_URL, TRENDS_URL, EastmoneyAdapter
from datahub.db import HubDB
from datahub.realtime_serve import fetch_intraday
from datahub.serving import ApiError, V1API
from datahub.timeutil import now_shanghai, yyyymmdd
class FakeEastmoney(EastmoneyAdapter):
def __init__(self) -> None:
super().__init__(timeout=2)
self.urls: list[str] = []
def _get_json(self, url, params, referer):
self.urls.append(f"{url}|{params.get('ndays')}")
if url == TRENDS_URL:
return {"data": {"name": "中国平安", "code": "601318", "preClose": 56.36, "trends": []}}
if url == HIS_TRENDS_URL:
return {
"data": {
"name": "中国平安",
"code": "601318",
"preClose": 55.8,
"trends": [
"2026-09-07 09:30,55.80,55.90,56.00,55.70,100,5580.00,55.900",
"2026-09-07 15:00,56.10,56.20,56.30,56.00,200,11240.00,56.150",
"2026-09-08 09:30,0,0,0,0,0,0.00,0",
],
}
}
raise AdapterError(f"unexpected url {url}")
class EastmoneyIntradayLookbackTests(unittest.TestCase):
def test_empty_today_uses_latest_available_session(self):
adapter = FakeEastmoney()
payload = adapter.fetch_intraday("601318.SH")
self.assertEqual(adapter.urls, [f"{TRENDS_URL}|1", f"{TRENDS_URL}|5", f"{HIS_TRENDS_URL}|5"])
self.assertEqual(payload["trade_date"], "2026-09-07")
self.assertEqual([point["time"] for point in payload["points"]], ["09:30", "15:00"])
self.assertEqual(payload["points"][0]["close"], 55.9)
def test_preferred_date_keeps_that_session(self):
adapter = FakeEastmoney()
payload = adapter.fetch_intraday("601318.SH", "20260907")
self.assertEqual(payload["trade_date"], "2026-09-07")
self.assertEqual(len(payload["points"]), 2)
class IntradayLkgTests(unittest.TestCase):
def setUp(self) -> None:
self.tmp = tempfile.TemporaryDirectory()
self.db = HubDB(Path(self.tmp.name) / "hub.db")
def tearDown(self) -> None:
self.tmp.cleanup()
def test_source_failure_returns_last_known_good(self):
from datahub.realtime_serve import _envelope, _write_cache
payload = _envelope(
{
"entity_type": "stock",
"ts_code": "601318.SH",
"trade_date": "2026-09-07",
"previous_close": 55.8,
"points": [{"date": "2026-09-07", "time": "09:30", "close": 55.9}],
},
{
"tier": "provisional",
"trade_date": "20260907",
"source": "eastmoney:trends2",
"stale": False,
},
)
_write_cache(self.db, "intraday:601318.SH:today", payload, 20, "eastmoney:trends2")
self.db.execute(
"UPDATE rt_cache SET expires_at = ? WHERE cache_key = ?",
("2000-01-01T00:00:00+08:00", "intraday:601318.SH:today"),
)
with patch("datahub.realtime_serve.EastmoneyAdapter") as mocked:
mocked.return_value.fetch_intraday.side_effect = AdapterError("down")
recovered = fetch_intraday(self.db, "601318.SH")
self.assertTrue(recovered["meta"]["stale"])
self.assertEqual(recovered["data"]["points"][0]["close"], 55.9)
def test_source_failure_without_lkg_raises(self):
with patch("datahub.realtime_serve.EastmoneyAdapter") as mocked:
mocked.return_value.fetch_intraday.side_effect = AdapterError("down")
with self.assertRaises(Exception) as ctx:
fetch_intraday(self.db, "000001.SZ")
self.assertIn("intraday unavailable", str(ctx.exception))
class ServingIntradayDateTests(unittest.TestCase):
def setUp(self) -> None:
self.tmp = tempfile.TemporaryDirectory()
self.db = HubDB(Path(self.tmp.name) / "hub.db")
self.api = V1API(self.db, pipeline=None, settings=None)
def tearDown(self) -> None:
self.tmp.cleanup()
def _assert_usable_intraday(self, payload: dict) -> None:
data = payload["data"]
points = [point for point in data.get("points") or [] if float(point.get("close") or 0) > 0]
self.assertGreaterEqual(len(points), 1)
self.assertTrue(str(data.get("trade_date") or ""))
self.assertFalse((payload.get("meta") or {}).get("stale"))
def test_serving_omitted_or_empty_date_uses_today_and_returns_points(self) -> None:
today = yyyymmdd(now_shanghai())
omitted = self.api.handle("/v1/intraday/points", {"code": ["601318"]})
empty = self.api.handle("/v1/intraday/points", {"code": ["601318"], "date": [""]})
explicit = self.api.handle("/v1/intraday/points", {"code": ["601318"], "date": [today]})
self._assert_usable_intraday(omitted)
self._assert_usable_intraday(empty)
self._assert_usable_intraday(explicit)
self.assertEqual(omitted["data"]["trade_date"], empty["data"]["trade_date"])
self.assertEqual(explicit["data"]["trade_date"], omitted["data"]["trade_date"])
def test_serving_normalizes_empty_date_to_today_and_keeps_history(self) -> None:
today = yyyymmdd(now_shanghai())
captured: list[str] = []
def fake_fetch(db, code, date=""):
captured.append(date)
return {
"schema_version": 1,
"data": {
"trade_date": f"{date[:4]}-{date[4:6]}-{date[6:8]}",
"points": [{"date": f"{date[:4]}-{date[4:6]}-{date[6:8]}", "time": "09:30", "close": 55.9}],
},
"meta": {"stale": False, "trade_date": date},
}
with patch("datahub.realtime_serve.fetch_intraday", side_effect=fake_fetch):
omitted = self.api.handle("/v1/intraday/points", {"code": ["601318"]})
empty = self.api.handle("/v1/intraday/points", {"code": ["601318"], "date": [" "]})
history = self.api.handle("/v1/intraday/points", {"code": ["601318"], "date": ["20260907"]})
self.assertEqual(captured, [today, today, "20260907"])
self.assertEqual(omitted["data"]["trade_date"], f"{today[:4]}-{today[4:6]}-{today[6:8]}")
self.assertEqual(empty["data"]["trade_date"], omitted["data"]["trade_date"])
self.assertEqual(history["data"]["trade_date"], "2026-09-07")
def test_serving_invalid_date_is_invalid_argument(self) -> None:
with self.assertRaises(ApiError) as ctx:
self.api.handle("/v1/intraday/points", {"code": ["601318"], "date": ["not-a-date"]})
self.assertEqual(ctx.exception.code, "INVALID_ARGUMENT")
self.assertIn("invalid trade_date", ctx.exception.message)
def test_serving_missing_code_is_invalid_argument(self) -> None:
with self.assertRaises(ApiError) as ctx:
self.api.handle("/v1/intraday/points", {"date": [yyyymmdd(now_shanghai())]})
self.assertEqual(ctx.exception.code, "INVALID_ARGUMENT")
self.assertIn("code is required", ctx.exception.message)
def test_serving_no_data_keeps_source_unavailable(self) -> None:
with patch("datahub.realtime_serve.EastmoneyAdapter") as mocked:
mocked.return_value.fetch_intraday.side_effect = AdapterError("No intraday chart data returned")
with self.assertRaises(ApiError) as ctx:
self.api.handle("/v1/intraday/points", {"code": ["000001"]})
self.assertEqual(ctx.exception.code, "SOURCE_UNAVAILABLE")
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()