From 7ad445bc9f5c7ad4103d5c4f70d7a1415a516e77 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=96=BD=E5=B7=A5=E5=91=98?= <4846a27a-ba87-49c9-939d-7275f63efcc7@agents.multica.local> Date: Thu, 27 Aug 2026 14:50:32 +0000 Subject: [PATCH] =?UTF-8?q?fix(HEL-190):=20=E6=8C=89=E7=9C=9F=E5=AE=9E?= =?UTF-8?q?=E4=BA=A4=E6=98=93=E6=97=A5=E5=8E=86=E8=A1=A5=E9=BD=90=E6=9C=80?= =?UTF-8?q?=E8=BF=9160=E6=97=A5=E5=BF=AB=E7=85=A7=EF=BC=8C=E4=BF=AE?= =?UTF-8?q?=E5=A4=8D=E6=96=AD=E6=A1=A3=E5=90=8E=E5=8F=AA=E6=98=BE=E7=A4=BA?= =?UTF-8?q?=E5=BD=93=E5=A4=A9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 保留连续性过滤,新增可审计补档工具与备份步骤;周末/节假日与真缺档分开处理,支持重复执行与部分失败续跑。 Co-authored-by: Cursor Co-authored-by: multica-agent --- backend/features/market/backfill_history.py | 202 +++++++++++++ backend/features/market/repository.py | 25 ++ backend/features/market/service.py | 250 +++++++++++++-- backend/features/system/routes.py | 10 +- config/architecture-inventory.json | 14 +- docs/maintenance/行情历史补档.md | 69 +++++ frontend/shared/admin.js | 9 +- tests/test_snapshot_backfill.py | 317 ++++++++++++++++++++ tools/README.md | 3 + tools/backfill_recent_snapshots.py | 106 +++++++ 10 files changed, 974 insertions(+), 31 deletions(-) create mode 100644 backend/features/market/backfill_history.py create mode 100644 docs/maintenance/行情历史补档.md create mode 100644 tests/test_snapshot_backfill.py create mode 100644 tools/backfill_recent_snapshots.py diff --git a/backend/features/market/backfill_history.py b/backend/features/market/backfill_history.py new file mode 100644 index 0000000..29f06b5 --- /dev/null +++ b/backend/features/market/backfill_history.py @@ -0,0 +1,202 @@ +"""Auditable recent-trading-day snapshot backfill helpers. + +Planning and backup stay free of provider imports so feature boundary tests remain green. +The service layer supplies open trading dates from the live calendar and executes sync. +""" + +from __future__ import annotations + +import sqlite3 +from datetime import date, datetime, timedelta +from pathlib import Path +from typing import Any, Iterable + + +MAX_RANGE_TRADING_DAYS = 15 +MAX_RECENT_TRADING_DAYS = 60 +DEFAULT_RECENT_TRADING_DAYS = 60 + +# Tables touched by a successful historical dashboard sync. User / token / model +# tables must never appear here. +SNAPSHOT_BACKFILL_WRITE_TABLES = frozenset( + { + "dashboard_snapshots", + "data_snapshots", + "sync_runs", + } +) + + +def clamp_recent_lookback(lookback: int) -> int: + value = int(lookback) + if value < 1: + raise ValueError("回补交易日数量至少为 1。") + if value > MAX_RECENT_TRADING_DAYS: + raise ValueError(f"单次最多回补最近 {MAX_RECENT_TRADING_DAYS} 个交易日。") + return value + + +def calendar_window_start(end_date: str, lookback: int) -> str: + """Natural-day lower bound large enough to cover lookback open sessions.""" + end = datetime.strptime(end_date, "%Y%m%d").date() + span = max(40, int(lookback * 2) + 20) + return (end - timedelta(days=span)).strftime("%Y%m%d") + + +def select_open_trade_dates( + calendar_rows: Iterable[dict[str, Any]], + end_date: str, + lookback: int, +) -> list[str]: + """Pick the last ``lookback`` open SSE sessions on or before ``end_date``.""" + lookback = clamp_recent_lookback(lookback) + end = normalize_compact_date(end_date) + open_dates = sorted( + { + normalize_compact_date(str(row.get("cal_date") or "")) + for row in calendar_rows + if int(row.get("is_open") or 0) == 1 and row.get("cal_date") + } + ) + open_dates = [item for item in open_dates if item <= end] + if not open_dates: + raise ValueError("交易日历未返回可用交易日,请检查行情 Token。") + return open_dates[-lookback:] + + +def select_open_trade_dates_in_range( + calendar_rows: Iterable[dict[str, Any]], + start_date: str, + end_date: str, + *, + maximum: int = MAX_RANGE_TRADING_DAYS, +) -> tuple[list[str], list[str]]: + """Return (open_dates, skipped_non_trading_days) inside an inclusive range.""" + start = normalize_compact_date(start_date) + end = normalize_compact_date(end_date) + if start > end: + raise ValueError("开始日期不能晚于结束日期。") + open_set = { + normalize_compact_date(str(row.get("cal_date") or "")) + for row in calendar_rows + if int(row.get("is_open") or 0) == 1 and row.get("cal_date") + } + open_dates: list[str] = [] + skipped: list[str] = [] + cursor = datetime.strptime(start, "%Y%m%d").date() + last = datetime.strptime(end, "%Y%m%d").date() + while cursor <= last: + compact = cursor.strftime("%Y%m%d") + if compact in open_set: + open_dates.append(compact) + else: + skipped.append(compact) + cursor += timedelta(days=1) + if len(open_dates) > maximum: + raise ValueError(f"单次最多回补 {maximum} 个交易日。") + return open_dates, skipped + + +def classify_snapshot_coverage( + trade_dates: list[str], + existing_dates: Iterable[str], +) -> dict[str, Any]: + present_set = { + normalize_compact_date(item) + for item in existing_dates + if item + } + present = [item for item in trade_dates if item in present_set] + missing = [item for item in trade_dates if item not in present_set] + return { + "trade_dates": list(trade_dates), + "present": present, + "missing": missing, + "present_count": len(present), + "missing_count": len(missing), + } + + +def create_sqlite_backup( + source_path: Path, + backup_dir: Path, + *, + label: str = "pre-backfill", + stamped_at: datetime | None = None, +) -> Path: + """Create a timestamped SQLite backup via the native backup API.""" + source = Path(source_path) + if not source.exists(): + raise FileNotFoundError(f"数据库不存在:{source}") + stamp = (stamped_at or datetime.now().astimezone()).strftime("%Y%m%d-%H%M%S") + safe_label = "".join(ch if ch.isalnum() or ch in "-_" else "-" for ch in label).strip("-") or "backup" + backup_dir = Path(backup_dir) + backup_dir.mkdir(parents=True, exist_ok=True) + target = backup_dir / f"review-{safe_label}-{stamp}.db" + source_conn = sqlite3.connect(f"file:{source}?mode=ro", uri=True) + try: + target_conn = sqlite3.connect(target) + try: + source_conn.backup(target_conn) + target_conn.commit() + finally: + target_conn.close() + finally: + source_conn.close() + return target + + +def display_date(compact: str) -> str: + value = normalize_compact_date(compact) + return f"{value[:4]}-{value[4:6]}-{value[6:8]}" + + +def normalize_compact_date(value: str) -> str: + compact = str(value or "").replace("-", "").strip() + if len(compact) != 8 or not compact.isdigit(): + raise ValueError("日期格式应为 YYYY-MM-DD。") + datetime.strptime(compact, "%Y%m%d") + return compact + + +def build_backfill_audit( + *, + mode: str, + end_date: str, + lookback: int | None, + coverage: dict[str, Any], + skipped_non_trading_days: list[str] | None = None, + backup_path: str | None = None, + dry_run: bool = False, + results: list[dict[str, Any]] | None = None, +) -> dict[str, Any]: + results = list(results or []) + succeeded = [row for row in results if row.get("status") == "success"] + skipped = [row for row in results if row.get("status") == "skipped"] + failed = [row for row in results if row.get("status") == "failed"] + return { + "ok": not failed, + "mode": mode, + "dry_run": dry_run, + "end_date": display_date(end_date), + "lookback": lookback, + "backup_path": backup_path, + "write_tables": sorted(SNAPSHOT_BACKFILL_WRITE_TABLES), + "trade_dates": [display_date(item) for item in coverage.get("trade_dates") or []], + "present": [display_date(item) for item in coverage.get("present") or []], + "missing": [display_date(item) for item in coverage.get("missing") or []], + "skipped_non_trading_days": [ + display_date(item) for item in (skipped_non_trading_days or []) + ], + "present_count": int(coverage.get("present_count") or 0), + "missing_count": int(coverage.get("missing_count") or 0), + "results": results, + "succeeded_count": len(succeeded), + "skipped_count": len(skipped), + "failed_count": len(failed), + "created_dates": [ + str(row.get("trade_date") or "") + for row in succeeded + if row.get("action") == "created" + ], + } diff --git a/backend/features/market/repository.py b/backend/features/market/repository.py index 7492e41..15c4828 100644 --- a/backend/features/market/repository.py +++ b/backend/features/market/repository.py @@ -227,6 +227,31 @@ class MarketRepositoryMixin: result.append(payload) return result + def list_snapshot_trade_dates( + self, + start_date: str = "", + end_date: str = "", + ) -> list[str]: + clauses: list[str] = [] + parameters: list[Any] = [] + if start_date: + clauses.append("trade_date >= ?") + parameters.append(start_date) + if end_date: + clauses.append("trade_date <= ?") + parameters.append(end_date) + where = f"WHERE {' AND '.join(clauses)}" if clauses else "" + with self.connect() as connection: + rows = connection.execute( + f""" + SELECT trade_date FROM dashboard_snapshots + {where} + ORDER BY trade_date + """, + parameters, + ).fetchall() + return [str(row["trade_date"]) for row in rows] + def start_sync(self, trade_date: str, source: str) -> int: started_at = datetime.now().astimezone().isoformat(timespec="seconds") with self.connect() as connection: diff --git a/backend/features/market/service.py b/backend/features/market/service.py index 935037e..545d867 100644 --- a/backend/features/market/service.py +++ b/backend/features/market/service.py @@ -3,9 +3,11 @@ from __future__ import annotations import copy import re from datetime import date, datetime, time as dt_time, timedelta +from pathlib import Path from typing import Any from backend.bootstrap.config import ( + DATA_DIR, normalize_date, tushare_code, validate_stock_code, @@ -13,6 +15,17 @@ from backend.bootstrap.config import ( ) from backend.data.providers.ifind_client import IfindError from backend.data.providers.tushare_client import TushareClient, TushareError +from backend.features.market.backfill_history import ( + DEFAULT_RECENT_TRADING_DAYS, + MAX_RANGE_TRADING_DAYS, + build_backfill_audit, + calendar_window_start, + classify_snapshot_coverage, + create_sqlite_backup, + display_date, + select_open_trade_dates, + select_open_trade_dates_in_range, +) from backend.features.market.charts import ChartDataError from backend.features.market.insights import MarketInsightsService from backend.features.sentiment.engine import SENTIMENT_ENGINE_VERSION @@ -890,31 +903,226 @@ class MarketServiceMixin: "intraday": intraday_points, } - def backfill(self, start_date: str, end_date: str) -> list[dict[str, Any]]: - start = datetime.strptime(normalize_date(start_date), "%Y%m%d").date() - end = datetime.strptime(normalize_date(end_date), "%Y%m%d").date() - if start > end: - raise ValueError("开始日期不能晚于结束日期。") - weekdays = [] - current = start - while current <= end: - if current.weekday() < 5: - weekdays.append(current) - current += timedelta(days=1) - if len(weekdays) > 15: - raise ValueError("单次最多回补 15 个工作日。") - results = [] - for day in weekdays: - dashboard = self.sync_dashboard(day.strftime("%Y%m%d")) + def backfill( + self, + start_date: str = "", + end_date: str = "", + *, + lookback: int | None = None, + dry_run: bool = False, + force: bool = False, + create_backup: bool = True, + ) -> dict[str, Any]: + """Backfill dashboard snapshots for real trading days only. + + - Date-range mode keeps the admin UI contract (max 15 open sessions). + - Recent mode fills the last N open sessions (default/max 60). + Weekends and holidays are reported as skipped non-trading days, not errors. + """ + if not self.configured: + raise ValueError("公共行情尚未配置,无法回补历史快照。") + normalized_end = normalize_date(end_date or date.today().isoformat()) + if lookback is not None or not (start_date and end_date): + target_lookback = ( + DEFAULT_RECENT_TRADING_DAYS if lookback is None else int(lookback) + ) + return self.backfill_recent_trading_days( + end_date=normalized_end, + lookback=target_lookback, + dry_run=dry_run, + force=force, + create_backup=create_backup, + ) + return self._backfill_date_range( + start_date=normalize_date(start_date), + end_date=normalized_end, + dry_run=dry_run, + force=force, + create_backup=create_backup, + ) + + def backfill_recent_trading_days( + self, + end_date: str = "", + lookback: int = DEFAULT_RECENT_TRADING_DAYS, + *, + dry_run: bool = False, + force: bool = False, + create_backup: bool = True, + ) -> dict[str, Any]: + normalized_end = normalize_date(end_date or date.today().isoformat()) + trade_dates = self._load_recent_open_trade_dates(normalized_end, lookback) + existing = self.database.list_snapshot_trade_dates( + trade_dates[0], trade_dates[-1] + ) + coverage = classify_snapshot_coverage(trade_dates, existing) + return self._execute_snapshot_backfill( + mode="recent", + end_date=normalized_end, + lookback=lookback, + coverage=coverage, + skipped_non_trading_days=[], + dry_run=dry_run, + force=force, + create_backup=create_backup, + ) + + def _backfill_date_range( + self, + start_date: str, + end_date: str, + *, + dry_run: bool = False, + force: bool = False, + create_backup: bool = True, + ) -> dict[str, Any]: + window_start = calendar_window_start(end_date, MAX_RANGE_TRADING_DAYS) + calendar_rows = self._tushare_client().query( + "trade_cal", + { + "exchange": "SSE", + "start_date": min(window_start, start_date), + "end_date": end_date, + }, + "cal_date,is_open,pretrade_date", + ) + trade_dates, skipped = select_open_trade_dates_in_range( + calendar_rows, + start_date, + end_date, + maximum=MAX_RANGE_TRADING_DAYS, + ) + if not trade_dates: + raise ValueError("选定区间内没有交易日,周末或节假日无需回补。") + existing = self.database.list_snapshot_trade_dates(trade_dates[0], trade_dates[-1]) + coverage = classify_snapshot_coverage(trade_dates, existing) + return self._execute_snapshot_backfill( + mode="range", + end_date=end_date, + lookback=None, + coverage=coverage, + skipped_non_trading_days=skipped, + dry_run=dry_run, + force=force, + create_backup=create_backup, + ) + + def _load_recent_open_trade_dates(self, end_date: str, lookback: int) -> list[str]: + start_date = calendar_window_start(end_date, lookback) + calendar_rows = self._tushare_client().query( + "trade_cal", + { + "exchange": "SSE", + "start_date": start_date, + "end_date": end_date, + }, + "cal_date,is_open,pretrade_date", + ) + return select_open_trade_dates(calendar_rows, end_date, lookback) + + def _execute_snapshot_backfill( + self, + *, + mode: str, + end_date: str, + lookback: int | None, + coverage: dict[str, Any], + skipped_non_trading_days: list[str], + dry_run: bool, + force: bool, + create_backup: bool, + ) -> dict[str, Any]: + targets = list(coverage["trade_dates"] if force else coverage["missing"]) + backup_path: str | None = None + if create_backup and not dry_run and targets: + backup = create_sqlite_backup( + Path(self.database.path), + DATA_DIR / "backups", + label=f"pre-{mode}-backfill", + ) + backup_path = str(backup) + + results: list[dict[str, Any]] = [] + if dry_run: + for trade_date in coverage["trade_dates"]: + exists = trade_date in coverage["present"] + if exists and not force: + status = "skipped" + action = "exists" + else: + status = "planned" + action = "refresh" if exists else "create" + results.append( + { + "requested_date": display_date(trade_date), + "trade_date": display_date(trade_date), + "status": status, + "action": action, + } + ) + return build_backfill_audit( + mode=mode, + end_date=end_date, + lookback=lookback, + coverage=coverage, + skipped_non_trading_days=skipped_non_trading_days, + backup_path=backup_path, + dry_run=True, + results=results, + ) + + present_before = set(coverage["present"]) + for trade_date in targets: + existed = trade_date in present_before + try: + dashboard = self.sync_dashboard(trade_date) + actual = normalize_date( + str(dashboard.get("meta", {}).get("trade_date") or trade_date) + ) + results.append( + { + "requested_date": display_date(trade_date), + "trade_date": display_date(actual), + "status": "success", + "action": "refreshed" if existed else "created", + "source": dashboard.get("meta", {}).get("source"), + "records": self._record_count(dashboard), + } + ) + except Exception as exc: + results.append( + { + "requested_date": display_date(trade_date), + "trade_date": display_date(trade_date), + "status": "failed", + "action": "refresh" if existed else "create", + "error": str(exc), + } + ) + + for trade_date in coverage["present"]: + if force: + continue results.append( { - "requested_date": day.isoformat(), - "trade_date": dashboard["meta"]["trade_date"], - "source": dashboard["meta"]["source"], - "records": self._record_count(dashboard), + "requested_date": display_date(trade_date), + "trade_date": display_date(trade_date), + "status": "skipped", + "action": "exists", } ) - return results + + results.sort(key=lambda row: str(row.get("requested_date") or "")) + return build_backfill_audit( + mode=mode, + end_date=end_date, + lookback=lookback, + coverage=coverage, + skipped_non_trading_days=skipped_non_trading_days, + backup_path=backup_path, + dry_run=False, + results=results, + ) def _stock_identity(self, code: str, trade_date: str) -> tuple[str, str]: snapshot = self.database.get_snapshot(trade_date) or {} diff --git a/backend/features/system/routes.py b/backend/features/system/routes.py index 5956cc9..1e0d840 100644 --- a/backend/features/system/routes.py +++ b/backend/features/system/routes.py @@ -29,11 +29,17 @@ class SystemRoutesMixin: def backfill_data(self) -> None: try: body = self.read_json_body() - results = self.application_service.backfill( + lookback_raw = body.get("lookback") + lookback = int(lookback_raw) if lookback_raw not in (None, "") else None + audit = self.application_service.backfill( str(body.get("start_date") or ""), str(body.get("end_date") or ""), + lookback=lookback, + dry_run=bool(body.get("dry_run")), + force=bool(body.get("force")), + create_backup=body.get("create_backup", True) is not False, ) - self.send_json({"ok": True, "results": results}) + self.send_json({"ok": True, **audit, "results": audit.get("results") or []}) except ValueError as exc: self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST) except Exception as exc: diff --git a/config/architecture-inventory.json b/config/architecture-inventory.json index e725ad6..ba3cc5d 100644 --- a/config/architecture-inventory.json +++ b/config/architecture-inventory.json @@ -541,8 +541,8 @@ }, { "path": "frontend/shared/admin.js", - "bytes": 14145, - "lines": 261 + "bytes": 14410, + "lines": 268 }, { "path": "backend/features/heaven/market_context.py", @@ -799,6 +799,11 @@ "bytes": 1919, "lines": 45 }, + { + "path": "backend/features/system/routes.py", + "bytes": 1791, + "lines": 46 + }, { "path": "backend/jobs/service.py", "bytes": 1746, @@ -829,11 +834,6 @@ "bytes": 1455, "lines": 48 }, - { - "path": "backend/features/system/routes.py", - "bytes": 1423, - "lines": 40 - }, { "path": "backend/features/themes/routes.py", "bytes": 1337, diff --git a/docs/maintenance/行情历史补档.md b/docs/maintenance/行情历史补档.md new file mode 100644 index 0000000..5425c56 --- /dev/null +++ b/docs/maintenance/行情历史补档.md @@ -0,0 +1,69 @@ +# 行情历史补档(最近 60 个交易日) + +用于修复 `dashboard_snapshots` 断档导致情绪周期 / 主题轮动 / 智能选股只剩当天的问题。 +保留 `latest_contiguous_history` 连续性规则;通过真实交易日历回补缺失交易日快照。 + +## 适用场景 + +- 库中已有稀疏历史快照,但最近一个真实交易日缺失,接口 `available_days=1`。 +- 需要可重复执行、可审计、可回退的补档,而不是迁库或放宽算法。 + +## 前置 + +1. 使用与线上一致的代码分支。 +2. 管理员账号已配置可用的公共 Tushare Token。 +3. 只操作目标环境自己的 `data/review.db`;禁止 `.36` 与 `.11` 互拷。 + +## 上线步骤(总工执行) + +在应用根目录: + +```bash +# 1) 只读规划:区分已有、真正缺档;不会写入 +python tools/backfill_recent_snapshots.py --account <管理员账号> --lookback 60 --dry-run --json + +# 2) 正式补档:先走 SQLite backup API 写 data/backups/review-pre-recent-backfill-*.db +# 再对缺失交易日调用现有 sync_dashboard +python tools/backfill_recent_snapshots.py --account <管理员账号> --lookback 60 --json + +# 3) 验证 +# GET /api/sentiment/history?trade_date=YYYY-MM-DD&limit=60 +# 期望 available_days >= 20,且不再只有 1 天 +``` + +管理端日期区间回补(`/api/backfill`)已改为只处理交易日历中的开市日,周末/节假日会进入 +`skipped_non_trading_days`,不再当成错误;单次仍限制 15 个交易日。最近 60 日请用本工具。 + +## 写入边界 + +只会通过现有同步路径写入: + +- `dashboard_snapshots` +- 同步审计表 `sync_runs` +- 必要时的 `data_snapshots`(仅当请求日被解析到其他交易日) + +不得改动用户、Token、模型绑定或系统配置表。 + +## 回滚 + +1. 优先按审计结果的 `created_dates` 精确删除新增行: + +```sql +DELETE FROM dashboard_snapshots WHERE trade_date IN ('YYYYMMDD', ...); +``` + +2. 若需整库回退,停止写入后用补档前备份覆盖: + +```bash +# 示例:把 data/backups/review-pre-recent-backfill-YYYYMMDD-HHMMSS.db +# 复制回 data/review.db 后重启容器 +``` + +3. 代码回退:对该提交执行 Git revert 后重新部署镜像。 + +## 验收要点 + +- dry-run 与正式执行可重复跑;已有交易日默认跳过。 +- 周末、节假日出现在 `skipped_non_trading_days`,不计入失败。 +- 部分交易日同步失败时,其他日期仍会继续,并在审计结果中标 `failed`。 +- 情绪周期、主题轮动 9 列、智能选股置信度随连续交易日恢复。 diff --git a/frontend/shared/admin.js b/frontend/shared/admin.js index 50dbbd2..cebe715 100644 --- a/frontend/shared/admin.js +++ b/frontend/shared/admin.js @@ -7,7 +7,14 @@ async function backfillData() { start_date: document.querySelector("#backfillStart").value, end_date: document.querySelector("#backfillEnd").value, }); - showToast(`历史回补完成,共处理 ${payload.results.length} 个工作日`); + const failed = (payload.failed_count || 0); + const skipped = (payload.skipped_non_trading_days || []).length; + const suffix = failed + ? `,失败 ${failed} 个` + : skipped + ? `,跳过 ${skipped} 个非交易日` + : ""; + showToast(`历史回补完成,共处理 ${payload.results.length} 个交易日${suffix}`); state.sentimentHistory = null; state.sentimentHistoryKey = ""; if (state.activeView === "sentimentCycleView") { diff --git a/tests/test_snapshot_backfill.py b/tests/test_snapshot_backfill.py new file mode 100644 index 0000000..a0b0b0e --- /dev/null +++ b/tests/test_snapshot_backfill.py @@ -0,0 +1,317 @@ +from __future__ import annotations + +import json +import tempfile +import threading +import unittest +from datetime import datetime +from pathlib import Path +from typing import Any +from unittest.mock import patch + +from backend.features.market.backfill_history import ( + build_backfill_audit, + classify_snapshot_coverage, + create_sqlite_backup, + select_open_trade_dates, + select_open_trade_dates_in_range, +) +from backend.features.market.service import MarketServiceMixin +from backend.features.sentiment.engine import ( + build_sentiment_history, + latest_contiguous_history, +) +from backend.features.sentiment.service import SentimentServiceMixin +from database import ReviewDatabase + + +def _snapshot(trade_date: str, previous_trade_date: str) -> dict[str, Any]: + display = f"{trade_date[:4]}-{trade_date[4:6]}-{trade_date[6:8]}" + previous_display = ( + f"{previous_trade_date[:4]}-{previous_trade_date[4:6]}-{previous_trade_date[6:8]}" + if previous_trade_date + else "" + ) + return { + "meta": { + "trade_date": display, + "previous_trade_date": previous_display, + "source": "tushare", + }, + "overview": { + "up_count": 2500, + "down_count": 2000, + "flat_count": 100, + "amount_billion": 12000, + "limit_up_count": 40, + "limit_down_count": 5, + "broken_count": 10, + "seal_rate": 70, + "max_height": 3, + "second_board_count": 8, + "three_plus_count": 4, + "previous_limit_count": 35, + "previous_positive_rate": 55, + "average_previous_change": 1.2, + "median_previous_change": 0.8, + "advance_rate": 20, + "severe_loss_rate": 5, + "previous_down_count": 3, + "ladder_completeness": 60, + "limit_amount_billion": 300, + }, + "limits": [{"code": "000001"}], + "broken": [], + "down_limits": [], + "yesterday_limits": [], + } + + +class _BackfillHarness(MarketServiceMixin, SentimentServiceMixin): + def __init__(self, database: ReviewDatabase) -> None: + self.database = database + self.sync_lock = threading.Lock() + self.configured = True + self.token = "test-token" + self.current_user_id = 1 + self._calendar_rows: list[dict[str, Any]] = [] + self._fail_dates: set[str] = set() + self.sync_calls: list[str] = [] + + def _tushare_client(self): # type: ignore[override] + harness = self + + class _Client: + def query(self, api_name, params, fields=""): + assert api_name == "trade_cal" + start = str(params["start_date"]) + end = str(params["end_date"]) + return [ + row + for row in harness._calendar_rows + if start <= str(row["cal_date"]) <= end + ] + + return _Client() + + def sync_dashboard(self, trade_date: str) -> dict[str, Any]: # type: ignore[override] + compact = trade_date.replace("-", "") + self.sync_calls.append(compact) + if compact in self._fail_dates: + raise ValueError(f"simulated failure for {compact}") + previous = "" + for row in self._calendar_rows: + if str(row["cal_date"]) == compact: + previous = str(row.get("pretrade_date") or "") + break + payload = _snapshot(compact, previous) + self.database.save_snapshot(compact, "tushare", payload) + return payload + + def _apply_reason_overrides(self, dashboard: dict[str, Any]) -> dict[str, Any]: + return dashboard + + def _with_storage(self, dashboard: dict[str, Any], cached: bool) -> dict[str, Any]: + return dashboard + + +class BackfillHistoryHelperTests(unittest.TestCase): + def test_select_open_trade_dates_skips_weekends_and_holidays(self) -> None: + rows = [ + {"cal_date": "20260821", "is_open": 1, "pretrade_date": "20260820"}, + {"cal_date": "20260822", "is_open": 0, "pretrade_date": "20260821"}, # Sat + {"cal_date": "20260823", "is_open": 0, "pretrade_date": "20260821"}, # Sun + {"cal_date": "20260824", "is_open": 1, "pretrade_date": "20260821"}, + {"cal_date": "20260825", "is_open": 1, "pretrade_date": "20260824"}, + {"cal_date": "20260826", "is_open": 1, "pretrade_date": "20260825"}, + {"cal_date": "20260827", "is_open": 1, "pretrade_date": "20260826"}, + ] + selected = select_open_trade_dates(rows, "20260827", 4) + self.assertEqual(selected, ["20260824", "20260825", "20260826", "20260827"]) + + def test_range_mode_reports_non_trading_days_separately(self) -> None: + rows = [ + {"cal_date": "20260821", "is_open": 1}, + {"cal_date": "20260824", "is_open": 1}, + ] + open_dates, skipped = select_open_trade_dates_in_range( + rows, "20260821", "20260824" + ) + self.assertEqual(open_dates, ["20260821", "20260824"]) + self.assertEqual(skipped, ["20260822", "20260823"]) + + def test_classify_snapshot_coverage_finds_real_gaps(self) -> None: + coverage = classify_snapshot_coverage( + ["20260824", "20260825", "20260826", "20260827"], + ["20260824", "20260827"], + ) + self.assertEqual(coverage["missing"], ["20260825", "20260826"]) + self.assertEqual(coverage["present"], ["20260824", "20260827"]) + + +class ContiguousHistoryGapTests(unittest.TestCase): + def test_missing_previous_trade_day_collapses_to_today(self) -> None: + payloads = [ + _snapshot("20260824", "20260821"), + _snapshot("20260827", "20260826"), # gap: 20260826 missing + ] + series = latest_contiguous_history(build_sentiment_history(payloads)) + self.assertEqual([row["trade_date"] for row in series], ["20260827"]) + + def test_continuous_history_keeps_full_tail(self) -> None: + payloads = [ + _snapshot("20260825", "20260824"), + _snapshot("20260826", "20260825"), + _snapshot("20260827", "20260826"), + ] + series = latest_contiguous_history(build_sentiment_history(payloads)) + self.assertEqual( + [row["trade_date"] for row in series], + ["20260825", "20260826", "20260827"], + ) + + +class SnapshotBackfillServiceTests(unittest.TestCase): + def setUp(self) -> None: + self.temporary = tempfile.TemporaryDirectory() + self.db_path = Path(self.temporary.name) / "review.db" + self.database = ReviewDatabase(self.db_path) + self.service = _BackfillHarness(self.database) + self.service._calendar_rows = [ + {"cal_date": "20260820", "is_open": 1, "pretrade_date": "20260819"}, + {"cal_date": "20260821", "is_open": 1, "pretrade_date": "20260820"}, + {"cal_date": "20260822", "is_open": 0, "pretrade_date": "20260821"}, + {"cal_date": "20260823", "is_open": 0, "pretrade_date": "20260821"}, + {"cal_date": "20260824", "is_open": 1, "pretrade_date": "20260821"}, + {"cal_date": "20260825", "is_open": 1, "pretrade_date": "20260824"}, + {"cal_date": "20260826", "is_open": 1, "pretrade_date": "20260825"}, + {"cal_date": "20260827", "is_open": 1, "pretrade_date": "20260826"}, + ] + # Sparse history mimicking .11: keep 0824 and today, miss 0825/0826. + self.database.save_snapshot("20260824", "tushare", _snapshot("20260824", "20260821")) + self.database.save_snapshot("20260827", "tushare", _snapshot("20260827", "20260826")) + + def tearDown(self) -> None: + self.temporary.cleanup() + + def test_recent_backfill_fills_gap_and_restores_history(self) -> None: + before = self.service.sentiment_history("20260827", 20) + self.assertEqual(before["available_days"], 1) + + with patch( + "backend.features.market.service.create_sqlite_backup", + return_value=Path(self.temporary.name) / "fake-backup.db", + ) as backup: + audit = self.service.backfill_recent_trading_days( + end_date="20260827", + lookback=4, + dry_run=False, + create_backup=True, + ) + + backup.assert_called_once() + self.assertEqual(sorted(self.service.sync_calls), ["20260825", "20260826"]) + self.assertEqual(audit["missing"], ["2026-08-25", "2026-08-26"]) + self.assertEqual(sorted(audit["created_dates"]), ["2026-08-25", "2026-08-26"]) + after = self.service.sentiment_history("20260827", 20) + self.assertGreaterEqual(after["available_days"], 4) + self.assertEqual( + [row["trade_date"] for row in after["rows"]], + ["20260824", "20260825", "20260826", "20260827"], + ) + + def test_dry_run_does_not_write_snapshots(self) -> None: + audit = self.service.backfill_recent_trading_days( + end_date="20260827", + lookback=4, + dry_run=True, + create_backup=True, + ) + self.assertTrue(audit["dry_run"]) + self.assertEqual(self.service.sync_calls, []) + self.assertIsNone(audit["backup_path"]) + self.assertEqual( + self.database.list_snapshot_trade_dates("20260824", "20260827"), + ["20260824", "20260827"], + ) + + def test_repeat_execution_skips_existing_days(self) -> None: + with patch( + "backend.features.market.service.create_sqlite_backup", + return_value=Path(self.temporary.name) / "fake-backup.db", + ): + first = self.service.backfill_recent_trading_days( + end_date="20260827", lookback=4 + ) + self.service.sync_calls.clear() + second = self.service.backfill_recent_trading_days( + end_date="20260827", lookback=4 + ) + self.assertEqual(first["succeeded_count"], 2) + self.assertEqual(self.service.sync_calls, []) + self.assertEqual(second["missing_count"], 0) + self.assertEqual(second["skipped_count"], 4) + self.assertIsNone(second["backup_path"]) + + def test_partial_failure_continues_remaining_days(self) -> None: + self.service._fail_dates.add("20260825") + with patch( + "backend.features.market.service.create_sqlite_backup", + return_value=Path(self.temporary.name) / "fake-backup.db", + ): + audit = self.service.backfill_recent_trading_days( + end_date="20260827", lookback=4 + ) + self.assertFalse(audit["ok"]) + self.assertEqual(audit["failed_count"], 1) + self.assertEqual(audit["succeeded_count"], 1) + self.assertIn("20260826", self.database.list_snapshot_trade_dates()) + self.assertNotIn("20260825", self.database.list_snapshot_trade_dates()) + + def test_range_backfill_skips_weekend_without_treating_as_error(self) -> None: + with patch( + "backend.features.market.service.create_sqlite_backup", + return_value=Path(self.temporary.name) / "fake-backup.db", + ): + audit = self.service.backfill( + start_date="2026-08-21", + end_date="2026-08-24", + ) + self.assertEqual(audit["mode"], "range") + self.assertEqual(audit["skipped_non_trading_days"], ["2026-08-22", "2026-08-23"]) + self.assertEqual(sorted(self.service.sync_calls), ["20260821"]) + self.assertTrue(audit["ok"]) + + def test_sqlite_backup_api_creates_restorable_copy(self) -> None: + backup_dir = Path(self.temporary.name) / "backups" + backup = create_sqlite_backup( + self.db_path, + backup_dir, + label="pre-recent-backfill", + stamped_at=datetime(2026, 8, 27, 15, 30, 0), + ) + self.assertTrue(backup.exists()) + self.assertIn("pre-recent-backfill-20260827-153000", backup.name) + restored = ReviewDatabase(backup) + self.assertEqual( + restored.list_snapshot_trade_dates(), + ["20260824", "20260827"], + ) + + def test_audit_lists_only_snapshot_related_write_tables(self) -> None: + audit = build_backfill_audit( + mode="recent", + end_date="20260827", + lookback=60, + coverage={"trade_dates": [], "present": [], "missing": [], "present_count": 0, "missing_count": 0}, + ) + self.assertEqual( + audit["write_tables"], + ["dashboard_snapshots", "data_snapshots", "sync_runs"], + ) + self.assertNotIn("users", audit["write_tables"]) + self.assertNotIn("system_settings", audit["write_tables"]) + + +if __name__ == "__main__": + unittest.main() diff --git a/tools/README.md b/tools/README.md index 417e3e4..04bece1 100644 --- a/tools/README.md +++ b/tools/README.md @@ -17,6 +17,9 @@ registry, and verification tools. `backend/features/*/routes.py` owners. - `python tools/build_architecture_inventory.py [--check]`: generate or verify `config/architecture-inventory.json` from the current source tree. +- `python tools/backfill_recent_snapshots.py --account [--lookback 60] [--dry-run]`: + auditable recent trading-day dashboard snapshot backfill. See + `docs/maintenance/行情历史补档.md`. `verify_baseline.py` does not inspect a parent checkout or skip tests according to files outside this application. Historical comparison scripts were retired after final standalone acceptance; diff --git a/tools/backfill_recent_snapshots.py b/tools/backfill_recent_snapshots.py new file mode 100644 index 0000000..df1e48e --- /dev/null +++ b/tools/backfill_recent_snapshots.py @@ -0,0 +1,106 @@ +#!/usr/bin/env python3 +"""Auditable recent trading-day dashboard snapshot backfill. + +Examples: + + python tools/backfill_recent_snapshots.py --account admin --dry-run + python tools/backfill_recent_snapshots.py --account admin --lookback 60 + python tools/backfill_recent_snapshots.py --account admin --end-date 2026-08-27 --force +""" + +from __future__ import annotations + +import argparse +import json +from datetime import date + +from backend.application import SERVICE +from backend.bootstrap.config import normalize_date +from backend.features.market.backfill_history import DEFAULT_RECENT_TRADING_DAYS + + +def main() -> None: + parser = argparse.ArgumentParser( + description="Backfill the latest N real trading-day dashboard snapshots" + ) + parser.add_argument( + "--account", + required=True, + help="Account that can resolve the shared Tushare token", + ) + parser.add_argument( + "--end-date", + default=date.today().isoformat(), + help="Inclusive end date YYYY-MM-DD (default: today)", + ) + parser.add_argument( + "--lookback", + type=int, + default=DEFAULT_RECENT_TRADING_DAYS, + help=f"Number of open trading days to cover (default {DEFAULT_RECENT_TRADING_DAYS}, max 60)", + ) + parser.add_argument( + "--dry-run", + action="store_true", + help="Plan only: classify missing gaps without writing", + ) + parser.add_argument( + "--force", + action="store_true", + help="Re-sync days that already have snapshots", + ) + parser.add_argument( + "--no-backup", + action="store_true", + help="Skip the SQLite backup API step (not recommended)", + ) + parser.add_argument( + "--json", + action="store_true", + help="Print the full audit payload as JSON", + ) + args = parser.parse_args() + + user = SERVICE.database.user_by_username(args.account.strip()) + if not user: + raise SystemExit("account not found") + SERVICE.bind_user(int(user["id"])) + + end_date = normalize_date(args.end_date) + audit = SERVICE.backfill_recent_trading_days( + end_date=end_date, + lookback=args.lookback, + dry_run=args.dry_run, + force=args.force, + create_backup=not args.no_backup, + ) + + if args.json: + print(json.dumps(audit, ensure_ascii=False, indent=2)) + raise SystemExit(0 if audit.get("ok") else 1) + + print( + f"mode={audit['mode']} end={audit['end_date']} lookback={audit['lookback']} " + f"dry_run={audit['dry_run']}" + ) + print( + f"present={audit['present_count']} missing={audit['missing_count']} " + f"succeeded={audit['succeeded_count']} skipped={audit['skipped_count']} " + f"failed={audit['failed_count']}" + ) + if audit.get("backup_path"): + print(f"backup={audit['backup_path']}") + if audit.get("missing"): + print("missing_dates=" + ",".join(audit["missing"])) + if audit.get("created_dates"): + print("created_dates=" + ",".join(audit["created_dates"])) + failed = [row for row in audit.get("results") or [] if row.get("status") == "failed"] + for row in failed: + print(f"failed {row.get('requested_date')}: {row.get('error')}") + if not audit.get("ok"): + raise SystemExit(1) + print("backfill complete") + + +if __name__ == "__main__": + main()