Compare commits

...
Author SHA1 Message Date
a043bc9eb1 fix(HEL-485): 盘中选择当天不再整页退回昨天
交易时段缺少盘后正式数据时继续展示当天盘中行情,只有开盘前、周末和历史日期才沿用最近收盘结果。

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
2026-09-08 10:16:52 +08:00
acde4de40d fix(HEL-484): 中枢分时接口空 date 按当天查询
缺少或为空的 date 不再 400,按当天处理;显式历史日期保持原行为。

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
2026-09-08 10:07:44 +08:00
3d2c1252f1 fix(HEL-482): 开盘前分时回退最近交易日,并接通中枢失败回旧通道
Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
2026-09-08 09:44:00 +08:00
18 changed files with 969 additions and 95 deletions
+48
View File
@@ -21,6 +21,26 @@ from backend.data.providers.tushare_client import TushareClient
LOGGER = logging.getLogger("xiaobai.datahub") LOGGER = logging.getLogger("xiaobai.datahub")
ShadowSink = Callable[[dict[str, Any]], None] 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 = { EMPTY_FAIL_DATASETS = {
"stocks", "daily", "index_daily", "valuation", "moneyflow", "auction", "stocks", "daily", "index_daily", "valuation", "moneyflow", "auction",
"limit_events", "sector_daily", "limit_events", "sector_daily",
@@ -91,6 +111,34 @@ class DatahubBridge:
self._log_failure("status", exc) self._log_failure("status", exc)
return None 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")
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 query( def query(
self, self,
api_name: str, api_name: str,
+3 -2
View File
@@ -85,12 +85,13 @@ def build_data_gateway(
policy = DataSourcePolicy.load() policy = DataSourcePolicy.load()
settings = datahub_settings or DatahubSettings.load(credentials=credentials) settings = datahub_settings or DatahubSettings.load(credentials=credentials)
datahub_client = DatahubClient(settings) datahub_client = DatahubClient(settings)
datahub = DatahubBridge(settings, datahub_client)
return DataGateway( return DataGateway(
policy=policy, policy=policy,
quality=DataQualityGate.load(policy), quality=DataQualityGate.load(policy),
tushare_provider=TushareProvider(token_supplier), tushare_provider=TushareProvider(token_supplier),
ifind_provider=IfindProvider(ifind), ifind_provider=IfindProvider(ifind),
chart_data=MarketChartClient(ifind, EastmoneyChartClient()), chart_data=MarketChartClient(ifind, EastmoneyChartClient(), datahub),
realtime_observer=WebRealtimeAggregator(), realtime_observer=WebRealtimeAggregator(),
datahub=DatahubBridge(settings, datahub_client), datahub=datahub,
) )
+10 -2
View File
@@ -3,7 +3,11 @@ from __future__ import annotations
from typing import Any from typing import Any
from backend.data.numbers import finite_number as _number from backend.data.numbers import finite_number as _number
from backend.data.providers.tushare_helpers import _display_time, _prices_equal from backend.data.providers.tushare_helpers import (
_display_time,
_prices_equal,
calendar_is_open,
)
class DailyMarketMixin: class DailyMarketMixin:
@@ -17,7 +21,11 @@ class DailyMarketMixin:
trade_date = requested trade_date = requested
else: else:
row = requested_rows[0] row = requested_rows[0]
trade_date = row["cal_date"] if row.get("is_open") == 1 else row.get("pretrade_date", requested) trade_date = (
row["cal_date"]
if calendar_is_open(row.get("is_open"))
else row.get("pretrade_date", requested)
)
resolved_rows = self.query( resolved_rows = self.query(
"trade_cal", "trade_cal",
+15 -9
View File
@@ -16,6 +16,12 @@ from backend.data.providers.tushare_transport import TushareError
class DashboardMixin: 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]: def dashboard(self, requested_date: str) -> dict[str, Any]:
trade_date, previous_trade_date = self.resolve_trade_context(requested_date) trade_date, previous_trade_date = self.resolve_trade_context(requested_date)
if self.should_use_realtime(requested_date, trade_date): if self.should_use_realtime(requested_date, trade_date):
@@ -26,11 +32,12 @@ class DashboardMixin:
) )
daily = self._load_daily(trade_date) daily = self._load_daily(trade_date)
now = self._now()
if ( if (
not daily not daily
and requested_date == datetime.now().astimezone().strftime("%Y%m%d") and requested_date == now.strftime("%Y%m%d")
and trade_date == requested_date and trade_date == requested_date
and datetime.now().astimezone().time().replace(tzinfo=None) >= dt_time(9, 15) and now.time().replace(tzinfo=None) >= dt_time(9, 15)
): ):
return self._realtime_dashboard( return self._realtime_dashboard(
requested_date, requested_date,
@@ -98,15 +105,14 @@ class DashboardMixin:
} }
return apply_sentiment_to_dashboard(dashboard) return apply_sentiment_to_dashboard(dashboard)
@staticmethod def should_use_realtime(self, requested_date: str, trade_date: str) -> bool:
def should_use_realtime(requested_date: str, trade_date: str) -> bool: """Use live quotes for today's open session until official daily settles."""
"""Use rt_k for today's open market until end-of-day datasets settle.""" now = self._now()
now = datetime.now().astimezone()
today = now.strftime("%Y%m%d") today = now.strftime("%Y%m%d")
return ( return (
requested_date == today requested_date == today
and trade_date == today and trade_date == today
and dt_time(9, 15) <= now.time().replace(tzinfo=None) < dt_time(16, 30) and dt_time(9, 15) <= now.time().replace(tzinfo=None) < dt_time(15, 5)
) )
def _realtime_dashboard( def _realtime_dashboard(
@@ -178,7 +184,7 @@ class DashboardMixin:
) )
sectors = _build_sectors(limits) sectors = _build_sectors(limits)
previous_sectors = _build_sectors(previous_limits) previous_sectors = _build_sectors(previous_limits)
now = datetime.now().astimezone() now = self._now()
market_status = _realtime_market_status(now.time().replace(tzinfo=None)) market_status = _realtime_market_status(now.time().replace(tzinfo=None))
dashboard = { dashboard = {
"meta": { "meta": {
@@ -234,7 +240,7 @@ class DashboardMixin:
{"trade_date": previous_trade_date}, {"trade_date": previous_trade_date},
"ts_code,trade_date,total_share,float_share,free_share,total_mv,circ_mv", "ts_code,trade_date,total_share,float_share,free_share,total_mv,circ_mv",
) )
if not basic_rows or not price_limits: if not basic_rows:
raise TushareError(f"Realtime reference data is incomplete for {trade_date}") raise TushareError(f"Realtime reference data is incomplete for {trade_date}")
result = { result = {
"basic_rows": basic_rows, "basic_rows": basic_rows,
+11
View File
@@ -6,6 +6,17 @@ from typing import Any
from backend.data.numbers import finite_number as _number 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: def _text(value: Any) -> str:
if isinstance(value, (list, tuple, set)): if isinstance(value, (list, tuple, set)):
return "".join(str(item).strip() for item in value if str(item).strip()) return "".join(str(item).strip() for item in value if str(item).strip())
+63 -15
View File
@@ -2,6 +2,7 @@ from __future__ import annotations
import http.client import http.client
import json import json
import logging
import re import re
import time import time
import urllib.error import urllib.error
@@ -15,12 +16,15 @@ from typing import Any, ClassVar
from backend.bootstrap.config import tushare_code as _stock_market_code from backend.bootstrap.config import tushare_code as _stock_market_code
from backend.data.providers.ifind_client import IfindError, IfindHttpClient from backend.data.providers.ifind_client import IfindError, IfindHttpClient
LOGGER = logging.getLogger("xiaobai.charts")
class ChartDataError(RuntimeError): class ChartDataError(RuntimeError):
pass pass
TRENDS_URL = "https://push2delay.eastmoney.com/api/qt/stock/trends2/get" 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" BOARD_LIST_URL = "https://push2delay.eastmoney.com/api/qt/clist/get"
BROWSER_USER_AGENT = ( BROWSER_USER_AGENT = (
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) " "Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
@@ -37,14 +41,23 @@ INDEX_SECIDS = {
class MarketChartClient: class MarketChartClient:
"""Prefer iFinD for display charts and retain Eastmoney as a last resort.""" """Prefer iFinD for display charts and retain Eastmoney as a last resort."""
def __init__(self, ifind: IfindHttpClient, fallback: "EastmoneyChartClient") -> None: def __init__(
self,
ifind: IfindHttpClient,
fallback: "EastmoneyChartClient",
datahub: Any = None,
) -> None:
self.ifind = ifind self.ifind = ifind
self.fallback = fallback self.fallback = fallback
self.datahub = datahub
def stock_intraday(self, code: str) -> dict[str, Any]: def stock_intraday(self, code: str) -> dict[str, Any]:
normalized = str(code or "").strip() normalized = str(code or "").strip()
if not re.fullmatch(r"\d{6}", normalized): if not re.fullmatch(r"\d{6}", normalized):
raise ChartDataError("Invalid stock code") 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) ifind_code = _stock_market_code(normalized)
try: try:
return self._ifind_intraday(ifind_code, "stock", normalized) return self._ifind_intraday(ifind_code, "stock", normalized)
@@ -73,11 +86,29 @@ class MarketChartClient:
normalized = str(identifier or "").strip().upper() normalized = str(identifier or "").strip().upper()
if normalized not in INDEX_SECIDS: if normalized not in INDEX_SECIDS:
raise ChartDataError("Unsupported index") raise ChartDataError("Unsupported index")
hub_chart = self._datahub_intraday(normalized)
if hub_chart is not None:
return hub_chart
try: try:
return self._ifind_intraday(normalized, "index", normalized) return self._ifind_intraday(normalized, "index", normalized)
except (IfindError, ChartDataError): except (IfindError, ChartDataError):
return self.fallback.index_intraday(normalized) 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 board_intraday(self, identifier: str, name: str = "") -> dict[str, Any]: def board_intraday(self, identifier: str, name: str = "") -> dict[str, Any]:
normalized = str(identifier or "").strip().upper() normalized = str(identifier or "").strip().upper()
try: try:
@@ -305,21 +336,29 @@ class EastmoneyChartClient:
if cached is not None: if cached is not None:
return cached return cached
payload = self._request_json( params = {
TRENDS_URL, "secid": secid,
{ "fields1": "f1,f2,f3,f4,f5,f6,f7,f8,f9,f10,f11,f12,f13",
"secid": secid, "fields2": "f51,f52,f53,f54,f55,f56,f57,f58",
"fields1": "f1,f2,f3,f4,f5,f6,f7,f8,f9,f10,f11,f12,f13", "iscr": "0",
"fields2": "f51,f52,f53,f54,f55,f56,f57,f58", }
"iscr": "0", last_error: Exception | None = None
"ndays": "1", data: dict[str, Any] = {}
}, points: list[dict[str, Any]] = []
"https://quote.eastmoney.com/", for url, ndays in ((TRENDS_URL, "1"), (TRENDS_URL, "5"), (HIS_TRENDS_URL, "5")):
) request_params = {**params, "ndays": ndays}
data = payload.get("data") or {} try:
points = [point for raw in data.get("trends") or [] if (point := _parse_trend(raw))] 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
if not points: if not points:
raise ChartDataError("No intraday chart data returned") raise ChartDataError("No intraday chart data returned") from last_error
result = { result = {
"entity_type": entity_type, "entity_type": entity_type,
@@ -433,6 +472,15 @@ class EastmoneyChartClient:
raise ChartDataError("Intraday chart request failed") from last_error 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: def _parse_trend(raw: Any) -> dict[str, Any] | None:
fields = str(raw or "").split(",") fields = str(raw or "").split(",")
if len(fields) < 8 or " " not in fields[0]: if len(fields) < 8 or " " not in fields[0]:
+64 -9
View File
@@ -65,9 +65,31 @@ class MarketServiceMixin:
# Compatibility for isolated legacy unit-test service stubs. # Compatibility for isolated legacy unit-test service stubs.
return TushareClient(self.token) return TushareClient(self.token)
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
def get_dashboard(self, trade_date: str, force: bool = False) -> dict[str, Any]: def get_dashboard(self, trade_date: str, force: bool = False) -> dict[str, Any]:
normalized_date = normalize_date(trade_date) normalized_date = normalize_date(trade_date)
now = datetime.now().astimezone() now = self._now()
if ( if (
normalized_date == now.strftime("%Y%m%d") normalized_date == now.strftime("%Y%m%d")
and now.time().replace(tzinfo=None) < datetime.strptime("09:15", "%H:%M").time() and now.time().replace(tzinfo=None) < datetime.strptime("09:15", "%H:%M").time()
@@ -174,14 +196,14 @@ class MarketServiceMixin:
def _should_retry_incomplete_snapshot( def _should_retry_incomplete_snapshot(
self, snapshot: dict[str, Any], requested_date: str self, snapshot: dict[str, Any], requested_date: str
) -> bool: ) -> bool:
if requested_date != date.today().strftime("%Y%m%d"): if requested_date != self._now().strftime("%Y%m%d"):
return False return False
meta = snapshot.get("meta") or {} meta = snapshot.get("meta") or {}
incomplete = ( actual = str(meta.get("trade_date") or "").replace("-", "")
meta.get("limit_data_source") == "derived" stale_carry = bool(meta.get("carried_forward") or actual != requested_date)
or bool(meta.get("carried_forward")) if stale_carry and self._is_requested_open_session(requested_date):
or str(meta.get("trade_date") or "").replace("-", "") != requested_date return True
) incomplete = meta.get("limit_data_source") == "derived" or stale_carry
return incomplete and self._snapshot_age_seconds(meta) >= 60 return incomplete and self._snapshot_age_seconds(meta) >= 60
def _annotate_data_status(self, dashboard: dict[str, Any]) -> dict[str, Any]: def _annotate_data_status(self, dashboard: dict[str, Any]) -> dict[str, Any]:
@@ -199,6 +221,9 @@ class MarketServiceMixin:
else: else:
meta["data_status"] = "preparing" meta["data_status"] = "preparing"
meta["display_notice"] = self._preparing_display_notice(actual, requested) meta["display_notice"] = self._preparing_display_notice(actual, requested)
elif meta.get("realtime"):
meta["data_status"] = "intraday"
meta.setdefault("display_notice", "")
else: else:
meta["data_status"] = "official" meta["data_status"] = "official"
meta.setdefault("display_notice", "") meta.setdefault("display_notice", "")
@@ -225,9 +250,9 @@ class MarketServiceMixin:
normalized_date: str, normalized_date: str,
snapshot: dict[str, Any], snapshot: dict[str, Any],
) -> bool: ) -> bool:
if not self.configured or normalized_date != date.today().strftime("%Y%m%d"): if not self.configured or normalized_date != self._now().strftime("%Y%m%d"):
return False return False
now = datetime.now().astimezone() now = self._now()
local_time = now.time().replace(tzinfo=None) local_time = now.time().replace(tzinfo=None)
realtime_start = datetime.strptime("09:15", "%H:%M").time() realtime_start = datetime.strptime("09:15", "%H:%M").time()
morning_end = datetime.strptime("11:35", "%H:%M").time() morning_end = datetime.strptime("11:35", "%H:%M").time()
@@ -276,6 +301,12 @@ class MarketServiceMixin:
actual_date = normalize_date( actual_date = normalize_date(
str(dashboard.get("meta", {}).get("trade_date") or normalized_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) self.database.save_snapshot(actual_date, source, dashboard)
if actual_date != normalized_date: if actual_date != normalized_date:
dashboard.setdefault("meta", {}).update( dashboard.setdefault("meta", {}).update(
@@ -297,6 +328,30 @@ class MarketServiceMixin:
) )
return self._apply_reason_overrides(self._with_storage(dashboard, cached=False)) return self._apply_reason_overrides(self._with_storage(dashboard, cached=False))
except TushareError as exc: 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) fallback = self.database.get_latest_real_snapshot(normalized_date)
if fallback: if fallback:
actual = str((fallback.get("meta") or {}).get("trade_date") or "") actual = str((fallback.get("meta") or {}).get("trade_date") or "")
+2
View File
@@ -41,6 +41,8 @@ def official_catchup_due(today: str, snapshot: dict[str, object]) -> bool:
actual == today actual == today
and meta.get("limit_data_source") != "derived" and meta.get("limit_data_source") != "derived"
and not meta.get("carried_forward") and not meta.get("carried_forward")
and not meta.get("realtime")
and meta.get("mode") != "realtime"
): ):
return False return False
return True return True
+19 -19
View File
@@ -508,8 +508,8 @@
}, },
{ {
"path": "backend/data/providers/tushare_dashboard.py", "path": "backend/data/providers/tushare_dashboard.py",
"bytes": 28234, "bytes": 28327,
"lines": 648 "lines": 654
}, },
{ {
"path": "backend/data/providers/tushare_industries.py", "path": "backend/data/providers/tushare_industries.py",
@@ -551,6 +551,11 @@
"bytes": 15311, "bytes": 15311,
"lines": 387 "lines": 387
}, },
{
"path": "frontend/shared/dashboard.js",
"bytes": 15063,
"lines": 321
},
{ {
"path": "frontend/pages/pools/page.html", "path": "frontend/pages/pools/page.html",
"bytes": 14942, "bytes": 14942,
@@ -561,11 +566,6 @@
"bytes": 14743, "bytes": 14743,
"lines": 342 "lines": 342
}, },
{
"path": "frontend/shared/dashboard.js",
"bytes": 14740,
"lines": 316
},
{ {
"path": "frontend/shared/admin.js", "path": "frontend/shared/admin.js",
"bytes": 14410, "bytes": 14410,
@@ -638,8 +638,8 @@
}, },
{ {
"path": "backend/data/providers/tushare_daily.py", "path": "backend/data/providers/tushare_daily.py",
"bytes": 6837, "bytes": 6949,
"lines": 160 "lines": 168
}, },
{ {
"path": "backend/application.py", "path": "backend/application.py",
@@ -786,6 +786,11 @@
"bytes": 2514, "bytes": 2514,
"lines": 63 "lines": 63
}, },
{
"path": "backend/data/providers/tushare_helpers.py",
"bytes": 2360,
"lines": 75
},
{ {
"path": "backend/jobs/service.py", "path": "backend/jobs/service.py",
"bytes": 2337, "bytes": 2337,
@@ -811,11 +816,6 @@
"bytes": 2165, "bytes": 2165,
"lines": 35 "lines": 35
}, },
{
"path": "backend/data/providers/tushare_helpers.py",
"bytes": 2083,
"lines": 64
},
{ {
"path": "frontend/pages/market/breadth.js", "path": "frontend/pages/market/breadth.js",
"bytes": 2071, "bytes": 2071,
@@ -827,13 +827,13 @@
"lines": 45 "lines": 45
}, },
{ {
"path": "backend/features/system/routes.py", "path": "backend/jobs/refresh.py",
"bytes": 1791, "bytes": 1808,
"lines": 46 "lines": 48
}, },
{ {
"path": "backend/jobs/refresh.py", "path": "backend/features/system/routes.py",
"bytes": 1728, "bytes": 1791,
"lines": 46 "lines": 46
}, },
{ {
+5
View File
@@ -67,6 +67,11 @@ async function startAdminRefresh() {
const actualCompact = actualDate.replaceAll("-", ""); const actualCompact = actualDate.replaceAll("-", "");
const updated = formatTimestamp(meta.updated_at); const updated = formatTimestamp(meta.updated_at);
const freshness = dashboardFreshnessMessage(meta); 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") { if (freshness || actualCompact !== requestedCompact || meta.carried_forward || meta.limit_data_source === "derived") {
setAdminRefreshStatus("warning", freshness || `部分正式数据尚未到齐,当前展示 ${actualDate || "最近可用数据"}`, "triangle-alert"); setAdminRefreshStatus("warning", freshness || `部分正式数据尚未到齐,当前展示 ${actualDate || "最近可用数据"}`, "triangle-alert");
setStatus(freshness || "部分正式数据尚未到齐,当前展示最近可用数据"); setStatus(freshness || "部分正式数据尚未到齐,当前展示最近可用数据");
+196 -15
View File
@@ -3,7 +3,8 @@ from __future__ import annotations
import copy import copy
import threading import threading
import unittest import unittest
from datetime import date, datetime, timedelta, timezone from datetime import date, datetime, timedelta, timezone, time as dt_time
from unittest.mock import patch
from pathlib import Path from pathlib import Path
from backend.features.market.service import MarketServiceMixin from backend.features.market.service import MarketServiceMixin
@@ -105,18 +106,58 @@ 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: class FakeMissingDailyClient:
def __init__(self, open_today: bool = True):
self.open_today = open_today
def dashboard(self, trade_date: str): def dashboard(self, trade_date: str):
raise TushareError(f"No daily data returned for {trade_date}") 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 SyncHarness(MarketServiceMixin): class SyncHarness(MarketServiceMixin):
def __init__(self, client, latest=None): def __init__(self, client, latest=None, clock=None):
self.configured = True self.configured = True
self.sync_lock = threading.Lock() self.sync_lock = threading.Lock()
self.database = FakeSyncDatabase(latest) self.database = FakeSyncDatabase(latest)
self._client = client self._client = client
self.current_user_id = 1 self.current_user_id = 1
self.clock = clock
def _tushare_client(self): def _tushare_client(self):
return self._client return self._client
@@ -142,23 +183,139 @@ class DashboardFreshnessTests(unittest.TestCase):
self.assertEqual(harness.database.finished[0][0][1], "success") self.assertEqual(harness.database.finished[0][0][1], "success")
self.assertEqual(verified_dashboard_result(payload), payload) self.assertEqual(verified_dashboard_result(payload), payload)
def test_missing_official_data_keeps_previous_day_with_preparing_notice(self): def test_intraday_refresh_keeps_today_and_does_not_fall_back_to_yesterday(self):
today = date.today() today = TRADE_DAY.strftime("%Y%m%d")
previous = (today - timedelta(days=1)).strftime("%Y-%m-%d")
latest = { latest = {
"meta": {"trade_date": previous, "source": "tushare"}, "meta": {"trade_date": "2026-09-07", "source": "tushare"},
"overview": {"limit_up_count": 20}, "overview": {"limit_up_count": 20},
} }
harness = SyncHarness(FakeMissingDailyClient(), latest) harness = SyncHarness(
payload = harness.sync_dashboard(today.strftime("%Y%m%d")) FakeRealtimeTodayClient(),
latest,
clock=lambda: at_clock(10, 5),
)
payload = harness.sync_dashboard(today)
meta = payload["meta"] meta = payload["meta"]
self.assertTrue(meta["carried_forward"]) self.assertFalse(meta.get("carried_forward"))
self.assertEqual(meta["data_status"], "preparing") self.assertTrue(meta["realtime"])
self.assertIn("今日数据正在准备,当前展示", meta["display_notice"]) self.assertEqual(meta["data_status"], "intraday")
self.assertIn("", meta["display_notice"]) self.assertEqual(str(meta["trade_date"]).replace("-", ""), today)
self.assertNotIn("No daily data", meta["display_notice"]) self.assertNotIn("今日数据正在准备", meta.get("display_notice") or "")
self.assertNotEqual(verified_dashboard_result(payload).get("status"), "failed") 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)
def test_weekend_carry_is_not_labeled_as_preparing(self): def test_weekend_carry_is_not_labeled_as_preparing(self):
snapshot = { snapshot = {
@@ -200,19 +357,43 @@ class DashboardFreshnessTests(unittest.TestCase):
{"meta": {"trade_date": iso, "limit_data_source": "derived"}}, {"meta": {"trade_date": iso, "limit_data_source": "derived"}},
) )
now = datetime.now().astimezone().time().replace(tzinfo=None) now = datetime.now().astimezone().time().replace(tzinfo=None)
if datetime.strptime("15:05", "%H:%M").time() <= now < datetime.strptime("22:00", "%H:%M").time(): if dt_time(15, 5) <= now < dt_time(22, 0):
self.assertFalse(due) self.assertFalse(due)
self.assertTrue(derived_due) self.assertTrue(derived_due)
else: else:
self.assertFalse(due) self.assertFalse(due)
self.assertFalse(derived_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): class FrontendRefreshCopyTests(unittest.TestCase):
def test_dashboard_script_distinguishes_partial_from_failure(self): def test_dashboard_script_distinguishes_partial_from_failure(self):
script = (Path(__file__).resolve().parents[1] / "frontend" / "shared" / "dashboard.js").read_text(encoding="utf-8") 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("部分正式数据尚未到齐", script)
self.assertIn("盘中行情", script)
self.assertIn("meta.realtime && actualCompact === requestedCompact", script)
self.assertIn('job.status === "failed"', script) self.assertIn('job.status === "failed"', script)
failed_block = script.split("if (job.status === \"failed\")", 1)[1].split("const query", 1)[0] failed_block = script.split("if (job.status === \"failed\")", 1)[1].split("const query", 1)[0]
self.assertIn("后台刷新失败", failed_block) self.assertIn("后台刷新失败", failed_block)
+137 -1
View File
@@ -2,7 +2,8 @@ from __future__ import annotations
import unittest import unittest
from backend.features.market.charts import ChartDataError, EastmoneyChartClient from backend.data.providers.ifind_client import IfindHttpClient
from backend.features.market.charts import ChartDataError, EastmoneyChartClient, HIS_TRENDS_URL, MarketChartClient, TRENDS_URL
from server import DashboardService from server import DashboardService
@@ -72,6 +73,141 @@ class ChartDataProviderTests(unittest.TestCase):
self.client.stock_intraday("abc") 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):
self.chart = chart
self.error = error
self.calls: list[str] = []
def try_intraday(self, code):
self.calls.append(code)
if self.error:
raise self.error
return self.chart
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)
class ChartServiceStub: class ChartServiceStub:
@staticmethod @staticmethod
def _payload(code: str, name: str): def _payload(code: str, name: str):
+65
View File
@@ -64,9 +64,11 @@ class FakeClient(DatahubClient):
meta={"tier": "official", "trade_date": "20240902", "stale": False, "staleness_seconds": 0}, meta={"tier": "official", "trade_date": "20240902", "stale": False, "staleness_seconds": 0},
) )
self.paths: list[str] = [] self.paths: list[str] = []
self.calls: list[tuple[str, dict[str, Any]]] = []
def get(self, path: str, params: dict[str, Any] | None = None) -> DatahubResponse: def get(self, path: str, params: dict[str, Any] | None = None) -> DatahubResponse:
self.paths.append(path) 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: if TOKEN in json.dumps(params or {}) or TOKEN in path:
raise AssertionError("token leaked into url") raise AssertionError("token leaked into url")
if self.error: if self.error:
@@ -349,6 +351,69 @@ class DatahubBridgeTests(unittest.TestCase):
self.assertEqual(rows[0]["amount"], 2000.0) self.assertEqual(rows[0]["amount"], 2000.0)
self.assertEqual(len(legacy.calls), 1) 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.assertFalse(DatahubSettings.load(environ={}, credentials={}).flags("intraday").read)
def test_features_do_not_import_datahub_client(self) -> None: def test_features_do_not_import_datahub_client(self) -> None:
violations = [] violations = []
for path in (ROOT / "backend" / "features").rglob("*.py"): for path in (ROOT / "backend" / "features").rglob("*.py"):
+54
View File
@@ -1,8 +1,10 @@
from __future__ import annotations from __future__ import annotations
import unittest import unittest
from datetime import datetime, timedelta, timezone
from backend.data.providers.tushare_client import TushareClient from backend.data.providers.tushare_client import TushareClient
from backend.data.providers.tushare_helpers import calendar_is_open
class FakeRealtimeClient(TushareClient): class FakeRealtimeClient(TushareClient):
@@ -130,6 +132,58 @@ class RealtimeDashboardTests(unittest.TestCase):
self.assertEqual(dashboard["meta"]["limit_data_source"], "derived") self.assertEqual(dashboard["meta"]["limit_data_source"], "derived")
self.assertIn("日线数据推算", dashboard["meta"]["notice"]) 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)
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()
+48 -20
View File
@@ -14,6 +14,7 @@ from datahub.numbers import finite_number, round4
EASTMONEY_INDEX_URL = "https://push2.eastmoney.com/api/qt/ulist.np/get" EASTMONEY_INDEX_URL = "https://push2.eastmoney.com/api/qt/ulist.np/get"
EASTMONEY_CLIST_URL = "https://push2.eastmoney.com/api/qt/clist/get" EASTMONEY_CLIST_URL = "https://push2.eastmoney.com/api/qt/clist/get"
TRENDS_URL = "https://push2delay.eastmoney.com/api/qt/stock/trends2/get" 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 = ( BROWSER_UA = (
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) " "Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
"AppleWebKit/537.36 (KHTML, like Gecko) Chrome/138.0.0.0 Safari/537.36" "AppleWebKit/537.36 (KHTML, like Gecko) Chrome/138.0.0.0 Safari/537.36"
@@ -166,7 +167,7 @@ class EastmoneyAdapter(MarketAdapter):
) )
return result return result
def fetch_intraday(self, ts_code: str) -> dict[str, Any]: def fetch_intraday(self, ts_code: str, date: str = "") -> dict[str, Any]:
code = str(ts_code or "").upper() code = str(ts_code or "").upper()
if code in INDEX_SECIDS: if code in INDEX_SECIDS:
secid = INDEX_SECIDS[code] secid = INDEX_SECIDS[code]
@@ -178,25 +179,32 @@ class EastmoneyAdapter(MarketAdapter):
secid = f"{market}.{symbol}" secid = f"{market}.{symbol}"
entity = "stock" entity = "stock"
identifier = symbol identifier = symbol
payload = self._get_json( params = {
TRENDS_URL, "secid": secid,
{ "fields1": "f1,f2,f3,f4,f5,f6,f7,f8,f9,f10,f11,f12,f13",
"secid": secid, "fields2": "f51,f52,f53,f54,f55,f56,f57,f58",
"fields1": "f1,f2,f3,f4,f5,f6,f7,f8,f9,f10,f11,f12,f13", "iscr": "0",
"fields2": "f51,f52,f53,f54,f55,f56,f57,f58", }
"iscr": "0", data: dict[str, Any] = {}
"ndays": "1", points: list[dict[str, Any]] = []
}, last_error: Exception | None = None
referer="https://quote.eastmoney.com/", for url, ndays in ((TRENDS_URL, "1"), (TRENDS_URL, "5"), (HIS_TRENDS_URL, "5")):
) try:
data = payload.get("data") or {} payload = self._get_json(
points = [] url,
for raw in data.get("trends") or []: {**params, "ndays": ndays},
point = _parse_trend(raw) referer="https://quote.eastmoney.com/",
if point: )
points.append(point) 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
if not points: if not points:
raise AdapterError("No intraday chart data returned") raise AdapterError("No intraday chart data returned") from last_error
return { return {
"entity_type": entity, "entity_type": entity,
"identifier": identifier, "identifier": identifier,
@@ -227,6 +235,23 @@ class EastmoneyAdapter(MarketAdapter):
raise AdapterError(f"eastmoney request failed: {exc}") from exc 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 _parse_trend(raw: Any) -> dict[str, Any] | None: def _parse_trend(raw: Any) -> dict[str, Any] | None:
text = str(raw or "") text = str(raw or "")
parts = text.split(",") parts = text.split(",")
@@ -237,11 +262,14 @@ def _parse_trend(raw: Any) -> dict[str, Any] | None:
when = datetime.strptime(stamp, "%Y-%m-%d %H:%M") when = datetime.strptime(stamp, "%Y-%m-%d %H:%M")
except ValueError: except ValueError:
return None return None
close = round4(finite_number(parts[2]))
if close <= 0:
return None
return { return {
"time": when.strftime("%H:%M"), "time": when.strftime("%H:%M"),
"date": when.strftime("%Y-%m-%d"), "date": when.strftime("%Y-%m-%d"),
"open": round4(finite_number(parts[1])), "open": round4(finite_number(parts[1])),
"close": round4(finite_number(parts[2])), "close": close,
"high": round4(finite_number(parts[3])), "high": round4(finite_number(parts[3])),
"low": round4(finite_number(parts[4])), "low": round4(finite_number(parts[4])),
"avg_price": round4(finite_number(parts[7] if len(parts) > 7 else parts[2])), "avg_price": round4(finite_number(parts[7] if len(parts) > 7 else parts[2])),
+47 -2
View File
@@ -14,6 +14,7 @@ from datahub.adapters.eastmoney import EastmoneyAdapter
from datahub.adapters.tencent import TencentAdapter from datahub.adapters.tencent import TencentAdapter
from datahub.codes import resolve_code from datahub.codes import resolve_code
from datahub.db import HubDB from datahub.db import HubDB
from datahub.governance.lkg import LastKnownGood
from datahub.timeutil import isoformat, now_shanghai, yyyymmdd from datahub.timeutil import isoformat, now_shanghai, yyyymmdd
QUOTE_TTL = 60 QUOTE_TTL = 60
@@ -108,10 +109,13 @@ def fetch_intraday(db: HubDB, code: str, date: str = "") -> dict[str, Any]:
return cached return cached
adapter = EastmoneyAdapter() adapter = EastmoneyAdapter()
try: try:
payload_data = adapter.fetch_intraday(ts_code) payload_data = adapter.fetch_intraday(ts_code, date)
source = "eastmoney:trends2" source = "eastmoney:trends2"
except Exception as exc: except Exception as exc:
raise RealtimeApiError("SOURCE_UNAVAILABLE", f"intraday unavailable: {exc}") from 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
payload = _envelope( payload = _envelope(
payload_data, payload_data,
{ {
@@ -127,6 +131,47 @@ def fetch_intraday(db: HubDB, code: str, date: str = "") -> dict[str, Any]:
return payload 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: def _guess_ts_code(code: str) -> str | None:
raw = str(code or "").strip().upper() raw = str(code or "").strip().upper()
if "." in raw: if "." in raw:
+6 -1
View File
@@ -269,8 +269,13 @@ class V1API:
code = str(q.get("code") or "").strip() code = str(q.get("code") or "").strip()
if not code: if not code:
raise ApiError("INVALID_ARGUMENT", "code is required") raise ApiError("INVALID_ARGUMENT", "code is required")
raw_date = str(q.get("date") or "").strip()
try: try:
return fetch_intraday(self.db, code, yyyymmdd(q.get("date") or "")) 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)
except RealtimeApiError as exc: except RealtimeApiError as exc:
raise ApiError(exc.code, exc.message) from exc raise ApiError(exc.code, exc.message) from exc
@@ -0,0 +1,176 @@
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)
if __name__ == "__main__":
unittest.main()