diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index b4d9e4f..809bc8d 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -54,7 +54,9 @@ background scheduler feature repository mixins; do not add feature queries to it. - `backend/jobs/` owns job definitions, locks, retries, idempotency, and persisted run state. `backend/jobs/service.py` is the application-facing owner of scheduler start/stop, manual - refresh submission, and periodic refresh coordination. + refresh submission, and periodic refresh coordination. `backend/jobs/refresh.py` owns + whether a dashboard payload is a usable refresh result versus a failed job, and whether + after-hours official catch-up is due. - `backend/llm/` owns model selection, membership/quota checks, fallback, provider transport, streaming rules, and call audit. Feature agents only prepare messages and interpret feature-specific results. diff --git a/README.md b/README.md index 368e786..4a192a8 100644 --- a/README.md +++ b/README.md @@ -16,7 +16,7 @@ - **题材库 / 人气热榜 / 龙虎榜**:题材成分、双榜人气、席位与游资档案 - **智能选股**(会员):六阶段策略、精选策略库、自然语言编译为受控公式后的确定性筛选与滚动回测;候选需手动加入后才进入五交易日跟踪 - **问师**(会员):按选定的游资思维 Skill 单师对话;新增公开角色时在 `游资skills` 下增加含 `SKILL.md` 的目录,并在 `游资skills/mentor_catalog.json` 登记。管理员私有角色放在 `data/private-mentor-skills`(不进 Git / 镜像) -- **问天**(会员,冻结区,勿改代码):观势 / 观气 / 观心。卦象、干支、节气与气机由本地程序确定性计算,大模型只负责文字解释 +- **问天**(会员):观势 / 观气 / 观心。卦象、干支、节气与气机由本地程序确定性计算,大模型只负责文字解释。此前仅冻结过界面视觉方案,现已解冻;问天可纳入后续数据与功能迁移,本阶段不主动重做视觉。 - **我的复盘**:手工交易日志、每日复盘、提醒中心与复盘助手;不接券商、不自动下单 全局能力:日间 / 夜间主题、股票代码悬停预览日 K 与分时、`Ctrl + K` 全局搜索。图表数据不写入主行情,也不参与情绪、选股或问天计算。 @@ -132,7 +132,7 @@ compose.yaml - 本项目是个人研究与复盘工具,全部数据、指标、候选与文字分析均不构成投资建议、证券推荐或买卖要约。 - 不接券商、不代为下单。交易日志只做手工记录与统计,不代表实际成交。 - 情绪温度、阶段判定、连板梯队、策略筛选等均为基于公开数据的统计与规则计算,不预测走势,不保证收益。 -- 「问天」属于传统文化视角的观察工具,不具备预测功能,不得作为投资依据;该模块为冻结区,不要改其代码。 +- 「问天」属于传统文化视角的观察工具,不具备预测功能,不得作为投资依据。问天不是永久冻结区:此前只冻结过界面视觉方案,现已解冻,后续数据与功能迁移可以纳入。 - 行情来自第三方接口,可能延迟、缺失或口径调整;不可用时页面会明确提示,请以交易所与券商正式披露为准。 - 不要把服务端口直接暴露到公网。不要把 Token、密码、密钥、数据库或 `.env` 提交进 Git。 - 股市有风险,入市需谨慎。投资决策及其后果由使用者本人承担。 diff --git a/backend/data/datahub/bridge.py b/backend/data/datahub/bridge.py index 068ca30..73ac792 100644 --- a/backend/data/datahub/bridge.py +++ b/backend/data/datahub/bridge.py @@ -25,6 +25,7 @@ EMPTY_FAIL_DATASETS = {"stocks", "daily", "index_daily", "valuation", "moneyflow def looks_like_heaven(module_name: str, filename: str = "") -> bool: + """问天调用栈识别。问天未永久冻结,只是本阶段仍走旧 Tushare 链路。""" path = filename.replace("\\", "/") return module_name.startswith("backend.features.heaven") or "/features/heaven/" in path @@ -95,6 +96,7 @@ class DatahubBridge: legacy_query: Callable[..., list[dict[str, Any]]], ) -> list[dict[str, Any]]: dataset = API_TO_DATASET.get(api_name) + # 问天允许后续纳入 datahub;首批只读接入仍保持旧链路,避免误切。 if not dataset or self.heaven_guard(): return legacy_query(api_name, params, fields) flags = self.settings.flags(dataset) diff --git a/backend/features/market/service.py b/backend/features/market/service.py index e372b8f..6680d68 100644 --- a/backend/features/market/service.py +++ b/backend/features/market/service.py @@ -79,6 +79,8 @@ class MarketServiceMixin: if not force: snapshot = self.database.get_snapshot(normalized_date) if snapshot and str((snapshot.get("meta") or {}).get("source") or "") != "demo": + if self._should_retry_incomplete_snapshot(snapshot, normalized_date): + return self.sync_dashboard(normalized_date) snapshot = copy.deepcopy(snapshot) if normalized_date != now.strftime("%Y%m%d"): snapshot.setdefault("meta", {}).update( @@ -97,6 +99,8 @@ class MarketServiceMixin: "dashboard_request_v1", normalized_date ) if resolved and str((resolved.get("meta") or {}).get("source") or "") != "demo": + if self._should_retry_incomplete_snapshot(resolved, normalized_date): + return self.sync_dashboard(normalized_date) resolved = copy.deepcopy(resolved) resolved.setdefault("meta", {})["requested_date"] = self._display_compact_date( normalized_date @@ -138,6 +142,68 @@ class MarketServiceMixin: def _display_compact_date(compact: str) -> str: return f"{compact[:4]}-{compact[4:6]}-{compact[6:8]}" + @staticmethod + def _chinese_month_day(value: str) -> str: + compact = str(value or "").replace("-", "").replace("/", "") + if len(compact) < 8 or not compact[:8].isdigit(): + return "最近可用交易日" + return f"{int(compact[4:6])} 月 {int(compact[6:8])} 日" + + @classmethod + def _preparing_display_notice(cls, actual_date: str, requested_date: str) -> str: + shown = cls._chinese_month_day(actual_date) + requested = str(requested_date or "").replace("-", "") + if requested == date.today().strftime("%Y%m%d"): + return f"今日数据正在准备,当前展示 {shown}" + return f"所选日期数据尚未到齐,当前展示 {shown}" + + @staticmethod + def _snapshot_age_seconds(meta: dict[str, Any]) -> float: + raw = str(meta.get("updated_at") or "") + if not raw: + return 10**9 + try: + updated_at = datetime.fromisoformat(raw) + except ValueError: + return 10**9 + now = datetime.now().astimezone() + if updated_at.tzinfo is None: + updated_at = updated_at.replace(tzinfo=now.tzinfo) + return (now - updated_at.astimezone(now.tzinfo)).total_seconds() + + def _should_retry_incomplete_snapshot( + self, snapshot: dict[str, Any], requested_date: str + ) -> bool: + if requested_date != date.today().strftime("%Y%m%d"): + return False + meta = snapshot.get("meta") or {} + 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]: + meta = dashboard.setdefault("meta", {}) + notice = str(meta.get("notice") or "") + requested = str(meta.get("requested_date") or "").replace("-", "") + actual = str(meta.get("trade_date") or "").replace("-", "") + if meta.get("limit_data_source") == "derived" and not meta.get("carried_forward"): + meta["data_status"] = "partial" + meta["display_notice"] = notice or "部分正式数据尚未到齐,当前展示日线推算结果" + elif meta.get("carried_forward"): + if "非交易日" in notice or "盘前" in notice: + meta["data_status"] = "carried" + meta["display_notice"] = notice + else: + meta["data_status"] = "preparing" + meta["display_notice"] = self._preparing_display_notice(actual, requested) + else: + meta["data_status"] = "official" + meta.setdefault("display_notice", "") + return dashboard + def _carry_dashboard( self, snapshot: dict[str, Any], requested_date: str, reason: str ) -> dict[str, Any]: @@ -152,7 +218,7 @@ class MarketServiceMixin: "notice": reason, } ) - return carried + return self._annotate_data_status(carried) def _realtime_snapshot_due( self, @@ -197,14 +263,14 @@ class MarketServiceMixin: if not self.configured: raise TushareError("公共行情尚未配置") dashboard = self._tushare_client().dashboard(normalized_date) - - if (dashboard.get("meta") or {}).get("limit_data_source") == "derived": - raise TushareError( - str((dashboard.get("meta") or {}).get("notice") or "官方涨跌停数据尚未返回") + meta = dashboard.setdefault("meta", {}) + meta["source"] = source + meta["requested_date"] = self._display_compact_date(normalized_date) + if meta.get("limit_data_source") == "derived": + meta.setdefault( + "notice", + "涨跌停高级接口当日数据尚未更新,已使用日线数据推算。", ) - - dashboard["meta"]["source"] = source - dashboard["meta"]["requested_date"] = self._display_compact_date(normalized_date) dashboard = self._enrich_dashboard_sentiment(dashboard, normalized_date) record_count = self._record_count(dashboard) actual_date = normalize_date( @@ -233,8 +299,11 @@ class MarketServiceMixin: except TushareError as exc: fallback = self.database.get_latest_real_snapshot(normalized_date) if fallback: + actual = str((fallback.get("meta") or {}).get("trade_date") or "") carried = self._carry_dashboard( - fallback, normalized_date, f"最新行情暂不可用,沿用最近收盘快照:{exc}" + fallback, + normalized_date, + self._preparing_display_notice(actual, normalized_date), ) self.database.finish_sync( sync_id, "fallback", self._record_count(carried), str(exc), "tushare" @@ -1160,7 +1229,7 @@ class MarketServiceMixin: "storage": "sqlite", "cached": cached, } - return result + return self._annotate_data_status(result) @staticmethod def _record_count(dashboard: dict[str, Any]) -> int: diff --git a/backend/http/handler.py b/backend/http/handler.py index 2a906fc..ed567c2 100644 --- a/backend/http/handler.py +++ b/backend/http/handler.py @@ -109,7 +109,14 @@ class HttpTransportMixin: return {} if length <= 0 or length > 65536: raise ValueError("请求内容为空或过大。") - return json.loads(self.rfile.read(length).decode("utf-8")) + raw = self.rfile.read(length) + try: + payload = json.loads(raw.decode("utf-8")) + except (UnicodeDecodeError, json.JSONDecodeError): + raise ValueError("请求不是合法 JSON。") from None + if not isinstance(payload, dict): + raise ValueError("请求不是合法 JSON。") + return payload def serve_static(self, request_path: str) -> None: relative = unquote(request_path).lstrip("/") or "index.html" diff --git a/backend/jobs/refresh.py b/backend/jobs/refresh.py new file mode 100644 index 0000000..279d9f8 --- /dev/null +++ b/backend/jobs/refresh.py @@ -0,0 +1,46 @@ +from __future__ import annotations + +from datetime import datetime, time as dt_time + + +def dashboard_has_usable_data(dashboard: dict[str, object]) -> bool: + if not isinstance(dashboard, dict) or dashboard.get("status") == "failed": + return False + meta = dashboard.get("meta") or {} + overview = dashboard.get("overview") or {} + if isinstance(meta, dict) and (meta.get("trade_date") or meta.get("carried_forward")): + return True + return bool(isinstance(overview, dict) and overview) + + +def verified_dashboard_result(dashboard: dict[str, object]) -> dict[str, object]: + """Manual refresh and automatic catch-up share this rule. + + Derived limit lists or a previous usable snapshot are not whole-job failures. + Only a payload with no displayable market data is recorded as failed. + """ + if dashboard_has_usable_data(dashboard): + return dashboard + meta = dashboard.get("meta") if isinstance(dashboard, dict) else None + notice = "" + if isinstance(meta, dict): + notice = str(meta.get("notice") or meta.get("display_notice") or "") + return { + "status": "failed", + "error": notice or "未获取到可用行情", + } + + +def official_catchup_due(today: str, snapshot: dict[str, object]) -> bool: + now = datetime.now().astimezone().time().replace(tzinfo=None) + if not (dt_time(15, 5) <= now < dt_time(22, 0)): + return False + meta = snapshot.get("meta") if isinstance(snapshot.get("meta"), dict) else {} + actual = str(meta.get("trade_date") or "").replace("-", "") + if ( + actual == today + and meta.get("limit_data_source") != "derived" + and not meta.get("carried_forward") + ): + return False + return True diff --git a/backend/jobs/service.py b/backend/jobs/service.py index 15dccc4..6be63cb 100644 --- a/backend/jobs/service.py +++ b/backend/jobs/service.py @@ -5,16 +5,7 @@ import time from datetime import date from backend.bootstrap.config import normalize_date - - -def _verified_dashboard_result(dashboard: dict[str, object]) -> dict[str, object]: - meta = dashboard.get("meta") or {} - if isinstance(meta, dict) and meta.get("carried_forward"): - return { - "status": "failed", - "error": str(meta.get("notice") or "未获取到所选日期的最新行情"), - } - return dashboard +from backend.jobs.refresh import official_catchup_due, verified_dashboard_result class JobServiceMixin: @@ -36,7 +27,7 @@ class JobServiceMixin: started = self.jobs.submit( "market.refresh", key, - lambda: _verified_dashboard_result(self.sync_dashboard(normalized)), + lambda: verified_dashboard_result(self.sync_dashboard(normalized)), {"trade_date": normalized, "trigger": "administrator"}, ) return {"started": started, "job_key": key if started else ""} @@ -54,7 +45,15 @@ class JobServiceMixin: self.jobs.submit( "market.refresh", f"realtime:{today}:{bucket}", - lambda: self.sync_dashboard(today), + lambda: verified_dashboard_result(self.sync_dashboard(today)), {"trade_date": today, "trigger": "realtime-poll"}, ) + elif official_catchup_due(today, snapshot): + bucket = int(time.time() // 300) + self.jobs.submit( + "market.refresh", + f"catchup:{today}:{bucket}", + lambda: verified_dashboard_result(self.sync_dashboard(today)), + {"trade_date": today, "trigger": "official-catchup"}, + ) self._schedule_automatic_screeners(today, snapshot) diff --git a/config/architecture-inventory.json b/config/architecture-inventory.json index 107c544..80a93e2 100644 --- a/config/architecture-inventory.json +++ b/config/architecture-inventory.json @@ -330,6 +330,7 @@ "system_service": "backend/features/system/service.py", "account_bridge": "backend/features/accounts/application.py", "job_lifecycle": "backend/jobs/service.py", + "job_refresh_status": "backend/jobs/refresh.py", "feature_routes": "backend/features/*/routes.py" }, "numeric_normalization": [ @@ -472,8 +473,8 @@ }, { "path": "frontend/shared/shell.css", - "bytes": 63659, - "lines": 3763 + "bytes": 63733, + "lines": 3767 }, { "path": "backend/features/heaven/engine.py", @@ -560,6 +561,11 @@ "bytes": 14743, "lines": 342 }, + { + "path": "frontend/shared/dashboard.js", + "bytes": 14740, + "lines": 316 + }, { "path": "frontend/shared/admin.js", "bytes": 14410, @@ -575,11 +581,6 @@ "bytes": 13219, "lines": 289 }, - { - "path": "frontend/shared/dashboard.js", - "bytes": 12894, - "lines": 274 - }, { "path": "backend/features/market/insights_auction_data.py", "bytes": 12829, @@ -785,16 +786,16 @@ "bytes": 2514, "lines": 63 }, + { + "path": "backend/jobs/service.py", + "bytes": 2337, + "lines": 59 + }, { "path": "backend/features/mentor/routes.py", "bytes": 2299, "lines": 57 }, - { - "path": "backend/jobs/service.py", - "bytes": 2219, - "lines": 60 - }, { "path": "backend/features/screener/regime.py", "bytes": 2202, @@ -830,6 +831,11 @@ "bytes": 1791, "lines": 46 }, + { + "path": "backend/jobs/refresh.py", + "bytes": 1728, + "lines": 46 + }, { "path": "backend/features/alerts/routes.py", "bytes": 1687, diff --git a/docs/README.md b/docs/README.md index d51cb26..6aed45c 100644 --- a/docs/README.md +++ b/docs/README.md @@ -32,4 +32,4 @@ - 旧文档不能删:被替代的旧文档开头要加一行「⚠️ 本文档已过时,仅留档备查,请勿删除」,再写新版。 - 用中文大白话写,专业词要带通俗解释,让不懂代码的人也能看懂。 -- 「问天」板块是冻结区,任何改动都不许碰;写文档时别误导后来人去改它。 +- 「问天」不是永久冻结区:此前只冻结过界面视觉方案,现已解冻。问天可纳入后续数据与功能迁移,不要再写成“永远不碰”。 diff --git a/docs/任务清单.md b/docs/任务清单.md index 815c771..6e587ef 100644 --- a/docs/任务清单.md +++ b/docs/任务清单.md @@ -7,6 +7,7 @@ | 任务 | 说明 | 状态 | |---|---|---| | 全站视觉统一改造收尾 | 主线。17 个阶段已完成,正在最终验收、代码合并 | 收尾中 | +| 行情刷新误报与旧数据提示 | HEL-412:高级接口未到齐不再记整次失败;今日正式数据晚到时提示当前展示日期 | 施工中 | | 手机端独立重新设计 | 先出视觉/交互规范和技术架构方案,等老板确认后再施工 | 方案送审中 | ## 已做完 diff --git a/docs/项目需求.md b/docs/项目需求.md index b01454a..5baf238 100644 --- a/docs/项目需求.md +++ b/docs/项目需求.md @@ -29,11 +29,11 @@ - **智能工具类(3 个)**:智能选股、问师、问天。 - **个人类(1 个)**:我的复盘。 -其中「问天」是冻结区(见下面的硬规矩)。 +其中「问天」此前只在全站视觉改造阶段冻结过界面方案,现已解冻;问天可以纳入后续数据与功能迁移,但不等于本阶段要重做视觉。 ## 几条硬规矩(不能破坏的边界) -- 「问天」板块是**冻结区**,任何改动都不许碰它。 +- 「问天」板块**不是永久冻结区**:此前冻结的是界面视觉方案,现已解冻。问天现有功能与界面不要破坏;后续数据与功能迁移可以纳入,不主动重做视觉。 - **不用假数据冒充真行情**;数据缺失就明说“没有/不可用”,不能编。 - **每个用户自己的数据互相隔离**(自选、复盘、对话、问天历史等),看不到别人的。 - **计算由程序确定性完成**(情绪周期、智能选股、问天排盘等),AI 大模型(LLM,就是会聊天的那个 AI)只负责解释或编译自然语言条件,不能改计算结果。 diff --git a/frontend/m/js/pages.js b/frontend/m/js/pages.js index 6f353dc..42fcab3 100644 --- a/frontend/m/js/pages.js +++ b/frontend/m/js/pages.js @@ -770,11 +770,33 @@ scroll.classList.add("m-motion-fade-in"); } + function dashboardFreshnessNotice() { + const meta = (state.dashboard && state.dashboard.meta) || {}; + if (meta.display_notice) return String(meta.display_notice); + const requested = String(meta.requested_date || "").replace(/-/g, ""); + const actual = String(meta.trade_date || "").replace(/-/g, ""); + const compact = actual; + const shown = /^\d{8}$/.test(compact) + ? (Number(compact.slice(4, 6)) + " 月 " + Number(compact.slice(6, 8)) + " 日") + : ""; + if (meta.data_status === "preparing" || (meta.carried_forward && actual && requested && actual !== requested)) { + return shown ? ("今日数据正在准备,当前展示 " + shown) : "今日数据正在准备,当前展示最近可用数据"; + } + if (meta.data_status === "partial" || meta.limit_data_source === "derived") { + return meta.notice || "部分正式数据尚未到齐,当前展示日线推算结果"; + } + return ""; + } + function renderTopArea(key) { const page = document.querySelector(".m-page"); if (!page) return; let top = page.querySelector(".m-top"); let html = buildStrip(); + const freshness = dashboardFreshnessNotice(); + if (freshness) { + html = '
' + escapeHtml(freshness) + "
" + html; + } if (key === "market/performance") html += performanceConclusion(); if (!top) { top = document.createElement("div"); diff --git a/frontend/shared/dashboard.js b/frontend/shared/dashboard.js index 54b612f..37c2d83 100644 --- a/frontend/shared/dashboard.js +++ b/frontend/shared/dashboard.js @@ -66,11 +66,13 @@ async function startAdminRefresh() { const requestedCompact = requestedDate.replaceAll("-", ""); const actualCompact = actualDate.replaceAll("-", ""); const updated = formatTimestamp(meta.updated_at); - if (actualCompact !== requestedCompact || meta.carried_forward) { - const reason = meta.notice ? `;${meta.notice}` : ""; - setAdminRefreshStatus("warning", `刷新已完成,但没有获取到 ${requestedDate} 的最新行情;当前仍是 ${actualDate || "未知日期"}${reason}`, "triangle-alert"); - showToast("刷新完成,但未获取到所选日期的最新行情"); - } else if (meta.notice) { + const freshness = dashboardFreshnessMessage(meta); + if (freshness || actualCompact !== requestedCompact || meta.carried_forward || meta.limit_data_source === "derived") { + setAdminRefreshStatus("warning", freshness || `部分正式数据尚未到齐,当前展示 ${actualDate || "最近可用数据"}`, "triangle-alert"); + setStatus(freshness || "部分正式数据尚未到齐,当前展示最近可用数据"); + return; + } + if (meta.notice) { setAdminRefreshStatus("warning", `已刷新到 ${actualDate}(${updated}),但数据源提示:${meta.notice}`, "triangle-alert"); showToast(`已刷新到 ${actualDate},请留意数据源提示`); } else { @@ -105,6 +107,37 @@ async function waitForAdminRefresh(jobKey) { throw new Error("刷新等待超时,请稍后重试"); } +let dashboardCatchupTimer = 0; + +function chineseMonthDay(value) { + const compact = String(value || "").replaceAll("-", "").replaceAll("/", ""); + if (!/^\d{8}/.test(compact)) return ""; + return `${Number(compact.slice(4, 6))} 月 ${Number(compact.slice(6, 8))} 日`; +} + +function dashboardFreshnessMessage(meta = {}) { + if (meta.display_notice) return String(meta.display_notice); + const requested = String(meta.requested_date || "").replaceAll("-", ""); + const actual = String(meta.trade_date || "").replaceAll("-", ""); + const shown = chineseMonthDay(actual); + if (meta.data_status === "preparing" || (meta.carried_forward && actual && requested && actual !== requested)) { + return shown ? `今日数据正在准备,当前展示 ${shown}` : "今日数据正在准备,当前展示最近可用数据"; + } + if (meta.data_status === "partial" || meta.limit_data_source === "derived") { + return meta.notice || "部分正式数据尚未到齐,当前展示日线推算结果"; + } + return ""; +} + +function scheduleDashboardCatchup(meta = {}) { + window.clearTimeout(dashboardCatchupTimer); + const status = String(meta.data_status || ""); + if (status !== "preparing" && status !== "partial") return; + dashboardCatchupTimer = window.setTimeout(() => { + loadDashboard(false, true, false); + }, 60000); +} + function applyDashboard(payload, background = false) { state.dashboard = payload; const selectedDate = payload.meta.requested_date || payload.meta.trade_date; @@ -112,7 +145,11 @@ function applyDashboard(payload, background = false) { document.querySelector("#qiObservationDate").value = selectedDate; document.querySelector("#journalDate").value = selectedDate; renderDashboard(); - setStatus(`${dashboardSourceLabel(payload.meta)} · 数据已更新`); + const freshness = dashboardFreshnessMessage(payload.meta || {}); + setStatus(freshness || `${dashboardSourceLabel(payload.meta)} · 数据已更新`); + const updatedAt = document.querySelector("#updatedAt"); + if (updatedAt) updatedAt.dataset.tone = freshness ? "warning" : "ok"; + scheduleDashboardCatchup(payload.meta || {}); if (!background) { if (state.activeView === "dragonView") loadDragonTiger(); if (state.activeView === "screenerView") loadScreenerSetup(); @@ -180,7 +217,12 @@ function renderDashboard() { } } updateSentimentGauge(overview.sentiment_score); - setText("updatedAt", `${dashboardSourceLabel(meta)} · 更新 ${formatTimestamp(meta.updated_at)}`); + const freshness = dashboardFreshnessMessage(meta); + setText("updatedAt", freshness + ? freshness + : `${dashboardSourceLabel(meta)} · 更新 ${formatTimestamp(meta.updated_at)}`); + const updatedAt = document.querySelector("#updatedAt"); + if (updatedAt) updatedAt.dataset.tone = freshness ? "warning" : "ok"; renderLimitTable(); renderLadderMini(ladders || []); diff --git a/frontend/shared/shell.css b/frontend/shared/shell.css index d76bd79..8daa720 100644 --- a/frontend/shared/shell.css +++ b/frontend/shared/shell.css @@ -921,6 +921,10 @@ body.sidebar-collapsed .app-main { text-align: right; } +.status-bar #updatedAt[data-tone="warning"] { + color: var(--warning); +} + .status-bar .risk-note { display: block; diff --git a/tests/test_admin_refresh_status.py b/tests/test_admin_refresh_status.py index 634f89e..ce65157 100644 --- a/tests/test_admin_refresh_status.py +++ b/tests/test_admin_refresh_status.py @@ -1,23 +1,223 @@ from __future__ import annotations +import copy +import threading import unittest +from datetime import date, datetime, timedelta, timezone +from pathlib import Path -from backend.jobs.service import _verified_dashboard_result +from backend.features.market.service import MarketServiceMixin +from backend.jobs.refresh import ( + dashboard_has_usable_data, + official_catchup_due, + verified_dashboard_result, +) +from backend.data.providers.tushare_transport import TushareError class AdminRefreshStatusTests(unittest.TestCase): - def test_carried_snapshot_is_reported_as_failed_job(self): - result = _verified_dashboard_result( - {"meta": {"carried_forward": True, "notice": "官方涨跌停数据尚未返回"}} + def test_carried_snapshot_is_usable_not_failed_job(self): + result = verified_dashboard_result( + { + "meta": { + "trade_date": "2026-09-01", + "requested_date": "2026-09-02", + "carried_forward": True, + "notice": "今日数据正在准备,当前展示 9 月 1 日", + "data_status": "preparing", + }, + "overview": {"limit_up_count": 12}, + } ) - self.assertEqual(result["status"], "failed") - self.assertEqual(result["error"], "官方涨跌停数据尚未返回") + self.assertNotEqual(result.get("status"), "failed") + self.assertEqual(result["meta"]["data_status"], "preparing") + self.assertTrue(dashboard_has_usable_data(result)) + + def test_derived_limit_snapshot_is_usable_not_failed_job(self): + dashboard = { + "meta": { + "trade_date": "2026-09-02", + "limit_data_source": "derived", + "notice": "涨跌停高级接口当日数据尚未更新,已使用日线数据推算。", + "data_status": "partial", + }, + "overview": {"limit_up_count": 8}, + } + + self.assertIs(verified_dashboard_result(dashboard), dashboard) def test_current_snapshot_is_reported_as_successful_job(self): dashboard = {"meta": {"trade_date": "2026-08-28", "carried_forward": False}} - self.assertIs(_verified_dashboard_result(dashboard), dashboard) + self.assertIs(verified_dashboard_result(dashboard), dashboard) + + def test_empty_payload_is_still_failed(self): + result = verified_dashboard_result({"meta": {}, "overview": {}}) + self.assertEqual(result["status"], "failed") + + +class FakeSyncDatabase: + def __init__(self, latest=None): + self.latest = latest + self.saved = [] + self.finished = [] + + def start_sync(self, *_args, **_kwargs): + return 1 + + def save_snapshot(self, trade_date, source, payload): + self.saved.append((trade_date, source, copy.deepcopy(payload))) + + def save_data_snapshot(self, *_args, **_kwargs): + return None + + def finish_sync(self, *args, **kwargs): + self.finished.append((args, kwargs)) + + def get_latest_real_snapshot(self, *_args, **_kwargs): + return copy.deepcopy(self.latest) + + def get_snapshot(self, *_args, **_kwargs): + return None + + def get_data_snapshot(self, *_args, **_kwargs): + return None + + def reason_overrides(self, *_args, **_kwargs): + return {} + + +class FakeDerivedClient: + def dashboard(self, trade_date: str): + return { + "meta": { + "trade_date": f"{trade_date[:4]}-{trade_date[4:6]}-{trade_date[6:8]}", + "limit_data_source": "derived", + "notice": "涨跌停高级接口当日数据尚未更新,已使用日线数据推算。", + "updated_at": datetime.now().astimezone().isoformat(timespec="seconds"), + }, + "overview": {"limit_up_count": 3}, + "limits": [{"code": "000001"}], + "broken": [], + "down_limits": [], + "yesterday_limits": [], + } + + +class FakeMissingDailyClient: + def dashboard(self, trade_date: str): + raise TushareError(f"No daily data returned for {trade_date}") + + +class SyncHarness(MarketServiceMixin): + 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 + + def _tushare_client(self): + return self._client + + def _enrich_dashboard_sentiment(self, dashboard, _trade_date): + return dashboard + + def _apply_reason_overrides(self, dashboard): + return dashboard + + +class DashboardFreshnessTests(unittest.TestCase): + def test_derived_limits_are_kept_as_partial_success(self): + today = date.today().strftime("%Y%m%d") + harness = SyncHarness(FakeDerivedClient()) + payload = harness.sync_dashboard(today) + meta = payload["meta"] + + self.assertEqual(meta["limit_data_source"], "derived") + self.assertEqual(meta["data_status"], "partial") + self.assertFalse(meta.get("carried_forward")) + self.assertIn("日线数据推算", meta["display_notice"]) + self.assertEqual(harness.database.finished[0][0][1], "success") + self.assertEqual(verified_dashboard_result(payload), payload) + + 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": previous, "source": "tushare"}, + "overview": {"limit_up_count": 20}, + } + harness = SyncHarness(FakeMissingDailyClient(), latest) + payload = harness.sync_dashboard(today.strftime("%Y%m%d")) + meta = payload["meta"] + + 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 = { + "meta": {"trade_date": "2026-07-24", "source": "tushare", "updated_at": "2026-07-24T15:00:00+08:00"}, + "overview": {"limit_up_count": 1}, + } + harness = SyncHarness(FakeMissingDailyClient()) + carried = harness._carry_dashboard(snapshot, "20260725", "非交易日沿用最近交易日收盘行情") + self.assertEqual(carried["meta"]["data_status"], "carried") + self.assertIn("非交易日", carried["meta"]["display_notice"]) + + def test_stale_derived_snapshot_is_retried(self): + today = date.today().strftime("%Y%m%d") + old = datetime.now(timezone.utc) - timedelta(minutes=5) + snapshot = { + "meta": { + "source": "tushare", + "trade_date": f"{today[:4]}-{today[4:6]}-{today[6:8]}", + "limit_data_source": "derived", + "updated_at": old.isoformat(), + }, + "overview": {"limit_up_count": 1}, + } + harness = SyncHarness(FakeDerivedClient()) + harness.database.get_snapshot = lambda *_args, **_kwargs: copy.deepcopy(snapshot) + payload = harness.get_dashboard(today) + self.assertEqual(payload["meta"]["data_status"], "partial") + self.assertTrue(harness.database.saved) + + def test_official_catchup_skips_complete_today_snapshot(self): + today = date.today().strftime("%Y%m%d") + iso = f"{today[:4]}-{today[4:6]}-{today[6:8]}" + due = official_catchup_due( + today, + {"meta": {"trade_date": iso, "limit_data_source": "official"}}, + ) + derived_due = official_catchup_due( + today, + {"meta": {"trade_date": iso, "limit_data_source": "derived"}}, + ) + now = datetime.now().astimezone().time().replace(tzinfo=None) + 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) + + +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('job.status === "failed"', script) + failed_block = script.split("if (job.status === \"failed\")", 1)[1].split("const query", 1)[0] + self.assertIn("后台刷新失败", failed_block) + success_block = script.split("const freshness = dashboardFreshnessMessage(meta);", 1)[1] + self.assertNotIn("后台刷新失败", success_block.split("} else {", 1)[0]) if __name__ == "__main__": diff --git a/tests/test_datahub_bridge.py b/tests/test_datahub_bridge.py index 56f1f75..c76abf8 100644 --- a/tests/test_datahub_bridge.py +++ b/tests/test_datahub_bridge.py @@ -198,7 +198,8 @@ class DatahubBridgeTests(unittest.TestCase): self.assertEqual(canonical["vol"], 100000.0) self.assertEqual(canonical["amount"], 2000000.0) - def test_heaven_keeps_legacy_even_when_read_flag_is_on(self) -> None: + def test_heaven_keeps_legacy_on_first_batch_even_when_read_flag_is_on(self) -> None: + """问天未永久冻结;首批只读接入仍走旧链路,后续迁移可以纳入。""" self.assertTrue(looks_like_heaven("backend.features.heaven.market_context", "backend/features/heaven/market_context.py")) self.assertFalse(looks_like_heaven("backend.features.market.service", "backend/features/market/service.py")) client = FakeClient() diff --git a/tools/build_architecture_inventory.py b/tools/build_architecture_inventory.py index 32c4bee..ddb12a3 100644 --- a/tools/build_architecture_inventory.py +++ b/tools/build_architecture_inventory.py @@ -97,6 +97,7 @@ def code_hotspots() -> list[dict[str, Any]]: "backend/features/system/service.py", "backend/features/accounts/application.py", "backend/jobs/service.py", + "backend/jobs/refresh.py", "database.py", "backend/features/screener/engine.py", "backend/features/screener/catalog.py", @@ -265,6 +266,7 @@ def build() -> dict[str, Any]: "system_service": "backend/features/system/service.py", "account_bridge": "backend/features/accounts/application.py", "job_lifecycle": "backend/jobs/service.py", + "job_refresh_status": "backend/jobs/refresh.py", "feature_routes": "backend/features/*/routes.py", }, "numeric_normalization": [ diff --git a/xiaobai-datahub/README.md b/xiaobai-datahub/README.md index bd59a44..207afe3 100644 --- a/xiaobai-datahub/README.md +++ b/xiaobai-datahub/README.md @@ -74,5 +74,6 @@ python -c "from pathlib import Path; from datahub.db import HubDB; HubDB(Path('d ## 安全 - 密钥只以 `configured / 末4位 / 更新时间` 出现在后台,不进日志、不进 `/v1` +- HTTP 解析失败只记录“请求不是合法 JSON”,不把请求正文、密码或 Token 写入容器日志 - 回滚、补数需重新输入密码 + 确认词 - 容器非 root(uid 10002)、read_only、cap_drop ALL diff --git a/xiaobai-datahub/datahub/adapters/tushare.py b/xiaobai-datahub/datahub/adapters/tushare.py index c16cd42..1b9b522 100644 --- a/xiaobai-datahub/datahub/adapters/tushare.py +++ b/xiaobai-datahub/datahub/adapters/tushare.py @@ -145,7 +145,9 @@ class TushareAdapter(MarketAdapter): try: with urllib.request.urlopen(request, timeout=self.timeout) as response: result = json.loads(response.read().decode("utf-8")) - except (urllib.error.URLError, TimeoutError, json.JSONDecodeError) as exc: + except json.JSONDecodeError: + raise AdapterError("Tushare returned invalid json") from None + except (urllib.error.URLError, TimeoutError) as exc: raise AdapterError(f"Tushare request failed: {exc}") from exc if result.get("code") != 0: raise AdapterError(result.get("msg") or "Tushare returned an unknown error") diff --git a/xiaobai-datahub/datahub/httpapp.py b/xiaobai-datahub/datahub/httpapp.py index de6328a..c649543 100644 --- a/xiaobai-datahub/datahub/httpapp.py +++ b/xiaobai-datahub/datahub/httpapp.py @@ -190,7 +190,15 @@ class HubRequestHandler(BaseHTTPRequestHandler): return {} if length <= 0 or length > 65536: raise ValueError("请求内容为空或过大") - return json.loads(self.rfile.read(length).decode("utf-8")) + raw = self.rfile.read(length) + try: + payload = json.loads(raw.decode("utf-8")) + except (UnicodeDecodeError, json.JSONDecodeError): + LOGGER.warning("invalid json request body") + raise ValueError("请求不是合法 JSON") from None + if not isinstance(payload, dict): + raise ValueError("请求不是合法 JSON") + return payload def _cookie_value(self, name: str) -> str: cookie = SimpleCookie() diff --git a/xiaobai-datahub/datahub/logutil.py b/xiaobai-datahub/datahub/logutil.py index 8752c62..f7f5851 100644 --- a/xiaobai-datahub/datahub/logutil.py +++ b/xiaobai-datahub/datahub/logutil.py @@ -2,7 +2,9 @@ from __future__ import annotations import json import logging +import re import sys +import traceback from typing import Any from datahub.timeutil import isoformat @@ -11,6 +13,13 @@ _SECRET_KEYS = ( "token", "password", "secret", "key", "authorization", "credential", "tushare_token", "datahub_token", "encryption_key", "cookie", ) +_SECRET_JSON = re.compile( + r'(?i)("(?:' + "|".join(re.escape(key) for key in _SECRET_KEYS) + r')"\s*:\s*")([^"\\]*(?:\\.[^"\\]*)*)(")' +) + + +def redact_log_text(text: str) -> str: + return _SECRET_JSON.sub(r"\1***\3", str(text)) def _redact(value: Any, key: str = "") -> Any: @@ -21,22 +30,39 @@ def _redact(value: Any, key: str = "") -> Any: return {str(item_key): _redact(item_value, str(item_key)) for item_key, item_value in value.items()} if isinstance(value, list): return [_redact(item) for item in value] + if isinstance(value, str): + return redact_log_text(value) return value +def _safe_exc_text(exc_info: tuple[Any, Any, Any]) -> str: + exc = exc_info[1] + if isinstance(exc, json.JSONDecodeError): + return f"JSONDecodeError: invalid json at position {exc.pos}" + cause = getattr(exc, "__cause__", None) + if isinstance(cause, json.JSONDecodeError): + return f"{type(exc).__name__}: invalid json in request" + text = "".join(traceback.format_exception(*exc_info)) + if isinstance(cause, json.JSONDecodeError) and cause.doc: + text = text.replace(cause.doc, "") + if isinstance(exc, json.JSONDecodeError) and exc.doc: + text = text.replace(exc.doc, "") + return redact_log_text(text) + + class JsonFormatter(logging.Formatter): def format(self, record: logging.LogRecord) -> str: payload: dict[str, Any] = { "ts": isoformat(), "level": record.levelname, "logger": record.name, - "message": record.getMessage(), + "message": redact_log_text(record.getMessage()), } extra = getattr(record, "hub", None) if isinstance(extra, dict): payload.update(_redact(extra)) if record.exc_info: - payload["exc"] = self.formatException(record.exc_info) + payload["exc"] = _safe_exc_text(record.exc_info) return json.dumps(payload, ensure_ascii=False, default=str) diff --git a/xiaobai-datahub/tests/test_admin.py b/xiaobai-datahub/tests/test_admin.py index 450411e..10a78db 100644 --- a/xiaobai-datahub/tests/test_admin.py +++ b/xiaobai-datahub/tests/test_admin.py @@ -1,17 +1,21 @@ from __future__ import annotations +import io import json +import logging import tempfile import threading import unittest from http.server import ThreadingHTTPServer from pathlib import Path +from urllib.error import HTTPError from urllib.request import Request, urlopen from datahub.adapters.tushare import TushareAdapter from datahub.crypto import SecretVault from datahub.httpapp import make_handler from datahub.hub import Hub +from datahub.logutil import JsonFormatter from datahub.settings import Settings from tests.fixtures import fake_transport @@ -91,6 +95,53 @@ class AdminTests(unittest.TestCase): ) self.assertEqual(ctx.exception.code, 401) + def test_invalid_json_does_not_log_request_body_secrets(self) -> None: + secret = "SuperSecretPass1!" + token = "hub-token-should-not-leak" + raw = json.dumps({"password": secret, "token": token, "username": "hub_admin"}) + "{not-json" + stream = io.StringIO() + logger = logging.getLogger("datahub") + handler = logging.StreamHandler(stream) + handler.setFormatter(JsonFormatter()) + logger.addHandler(handler) + previous_level = logger.level + logger.setLevel(logging.DEBUG) + try: + req = Request( + self.base + "/admin/api/login", + data=raw.encode("utf-8"), + headers={"Content-Type": "application/json"}, + method="POST", + ) + with self.assertRaises(HTTPError) as ctx: + urlopen(req, timeout=5) + body = ctx.exception.read().decode("utf-8") + self.assertEqual(ctx.exception.code, 400) + self.assertNotIn(secret, body) + self.assertNotIn(token, body) + blob = stream.getvalue() + body + self.assertNotIn(secret, blob) + self.assertNotIn(token, blob) + self.assertNotIn(raw, blob) + finally: + logger.removeHandler(handler) + logger.setLevel(previous_level) + + def test_json_formatter_drops_decode_error_document(self) -> None: + secret = "ParseSecretTokenXYZ" + formatter = JsonFormatter() + logger = logging.getLogger("datahub.test") + record = logger.makeRecord( + "datahub.test", logging.ERROR, __file__, 1, "parse failed", (), None + ) + try: + json.loads('{"password": "%s"}{' % secret) + except json.JSONDecodeError as exc: + record.exc_info = (type(exc), exc, exc.__traceback__) + blob = formatter.format(record) + self.assertNotIn(secret, blob) + self.assertIn("invalid json", blob) + if __name__ == "__main__": unittest.main()