Compare commits

...
Author SHA1 Message Date
总工andmultica-agent 6d7a839202 chore(HEL-208): 刷新 architecture-inventory 以对齐 HEL-207 行情改动
Co-authored-by: multica-agent <github@multica.ai>
2026-08-28 08:09:49 +00:00
总工 fc1e5b89e4 merge(HEL-208): 集成收盘行情修复 cf206c7 到 .11 正式线(基于 09a935a) 2026-08-28 08:09:28 +00:00
cf206c7de9 fix(HEL-207): 收盘后改走日线,禁止误调 rt_k,回退旧快照记失败
将实时窗口与调度窗口统一到 15:05;盘后优先日线并补关闭盘后同步;沿用旧快照时 sync/job 记 failed。

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
2026-08-28 08:00:08 +00:00
总工andmultica-agent 09a935aac4 merge(HEL-193): 以 .11 线上基线 013ed29 整合行情历史补档 8e94c7b(HEL-190/HEL-199)
- 合入 agent/agent/ebb4e3e0638a:真实交易日历补最近60日快照 + sys.path 引导返工
- architecture-inventory 于合并后重新生成

Co-authored-by: multica-agent <github@multica.ai>
2026-08-27 16:13:34 +00:00
8e94c7b429 fix(HEL-199): 为补档工具补上仓库根 sys.path 引导
使 python3 tools/backfill_recent_snapshots.py --help 在干净环境下可直接运行,并同步文档运行示例为容器内执行。

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
2026-08-27 15:12:42 +00:00
7ad445bc9f fix(HEL-190): 按真实交易日历补齐最近60日快照,修复断档后只显示当天
保留连续性过滤,新增可审计补档工具与备份步骤;周末/节假日与真缺档分开处理,支持重复执行与部分失败续跑。

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
2026-08-27 14:50:32 +00:00
13 changed files with 1255 additions and 51 deletions
+3 -13
View File
@@ -26,18 +26,8 @@ class DashboardMixin:
)
daily = self._load_daily(trade_date)
if (
not daily
and requested_date == datetime.now().astimezone().strftime("%Y%m%d")
and trade_date == requested_date
and datetime.now().astimezone().time().replace(tzinfo=None) >= dt_time(9, 15)
):
return self._realtime_dashboard(
requested_date,
trade_date,
previous_trade_date,
)
if not daily:
# 15:05 后只走日线;日线未就绪时不得回退调用无权限的 rt_k。
raise TushareError(f"No daily data returned for {trade_date}")
notices: list[str] = []
@@ -96,13 +86,13 @@ class DashboardMixin:
@staticmethod
def should_use_realtime(requested_date: str, trade_date: str) -> bool:
"""Use rt_k for today's open market until end-of-day datasets settle."""
"""Use rt_k only inside the intraday window; 15:05+ must use daily bars."""
now = datetime.now().astimezone()
today = now.strftime("%Y%m%d")
return (
requested_date == today
and trade_date == today
and dt_time(9, 15) <= now.time().replace(tzinfo=None) < dt_time(16, 30)
and dt_time(9, 15) <= now.time().replace(tzinfo=None) < dt_time(15, 5)
)
def _realtime_dashboard(
+202
View File
@@ -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"
],
}
+25
View File
@@ -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:
+261 -23
View File
@@ -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
@@ -175,6 +188,30 @@ class MarketServiceMixin:
age_seconds = (now - updated_at.astimezone(now.tzinfo)).total_seconds()
return age_seconds >= 8
def _closing_snapshot_due(
self,
normalized_date: str,
snapshot: dict[str, Any],
) -> bool:
"""After 15:05, keep requesting daily bars until today's EOD snapshot exists."""
if not self.configured or normalized_date != date.today().strftime("%Y%m%d"):
return False
now = datetime.now().astimezone()
if now.weekday() >= 5:
return False
local_time = now.time().replace(tzinfo=None)
if local_time < datetime.strptime("15:05", "%H:%M").time():
return False
meta = snapshot.get("meta") or {}
snapshot_trade_date = str(meta.get("trade_date") or "").replace("-", "")
if (
snapshot_trade_date == normalized_date
and not meta.get("realtime")
and not meta.get("carried_forward")
):
return False
return True
def sync_dashboard(self, trade_date: str) -> dict[str, Any]:
normalized_date = normalize_date(trade_date)
source = "tushare"
@@ -219,9 +256,15 @@ class MarketServiceMixin:
fallback, normalized_date, f"最新行情暂不可用,沿用最近收盘快照:{exc}"
)
self.database.finish_sync(
sync_id, "fallback", self._record_count(carried), str(exc), "tushare"
sync_id, "failed", self._record_count(carried), str(exc), "tushare"
)
return self._apply_reason_overrides(self._with_storage(carried, cached=True))
result = self._apply_reason_overrides(
self._with_storage(carried, cached=True)
)
# 页面仍可读到沿用快照;后台任务通过顶层 status=failed 记失败。
result["status"] = "failed"
result["error"] = str(exc)
return result
self.database.finish_sync(sync_id, "failed", message=str(exc))
raise ValueError("暂无可用的真实行情快照,请等待后台完成首次同步。") from exc
except Exception as exc:
@@ -890,31 +933,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 {}
+8 -2
View File
@@ -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:
+9
View File
@@ -46,4 +46,13 @@ class JobServiceMixin:
lambda: self.sync_dashboard(today),
{"trade_date": today, "trigger": "realtime-poll"},
)
elif self._closing_snapshot_due(today, snapshot):
# 15:05 后改走日线生成当日快照;按分钟去重,避免日线未就绪时刷爆任务。
bucket = int(time.time() // 60)
self.jobs.submit(
"market.refresh",
f"closing:{today}:{bucket}",
lambda: self.sync_dashboard(today),
{"trade_date": today, "trigger": "post-close"},
)
self._schedule_automatic_screeners(today, snapshot)
+12 -12
View File
@@ -486,8 +486,8 @@
},
{
"path": "backend/data/providers/tushare_dashboard.py",
"bytes": 28051,
"lines": 644
"bytes": 27730,
"lines": 634
},
{
"path": "backend/data/providers/tushare_industries.py",
@@ -541,8 +541,8 @@
},
{
"path": "frontend/shared/admin.js",
"bytes": 14145,
"lines": 261
"bytes": 14410,
"lines": 268
},
{
"path": "backend/features/heaven/market_context.py",
@@ -774,6 +774,11 @@
"bytes": 2202,
"lines": 53
},
{
"path": "backend/jobs/service.py",
"bytes": 2201,
"lines": 58
},
{
"path": "backend/data/providers/tushare_client.py",
"bytes": 2166,
@@ -800,9 +805,9 @@
"lines": 45
},
{
"path": "backend/jobs/service.py",
"bytes": 1746,
"lines": 49
"path": "backend/features/system/routes.py",
"bytes": 1791,
"lines": 46
},
{
"path": "backend/features/alerts/routes.py",
@@ -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,
+69
View File
@@ -0,0 +1,69 @@
# 行情历史补档(最近 60 个交易日)
用于修复 `dashboard_snapshots` 断档导致情绪周期 / 主题轮动 / 智能选股只剩当天的问题。
保留 `latest_contiguous_history` 连续性规则;通过真实交易日历回补缺失交易日快照。
## 适用场景
- 库中已有稀疏历史快照,但最近一个真实交易日缺失,接口 `available_days=1`
- 需要可重复执行、可审计、可回退的补档,而不是迁库或放宽算法。
## 前置
1. 使用与线上一致的代码分支。
2. 管理员账号已配置可用的公共 Tushare Token。
3. 只操作目标环境自己的 `data/review.db`;禁止 `.36``.11` 互拷。
## 上线步骤(总工执行)
在目标环境容器内执行(应用根目录;宿主机也可直接跑,脚本已自带仓库根 `sys.path` 引导):
```bash
# 1) 只读规划:区分已有、真正缺档;不会写入
docker compose exec xiaobai-review 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
docker compose exec xiaobai-review 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 列、智能选股置信度随连续交易日恢复。
+8 -1
View File
@@ -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") {
+226
View File
@@ -0,0 +1,226 @@
from __future__ import annotations
import threading
import unittest
from datetime import datetime
from unittest.mock import patch
from backend.data.providers.tushare_client import TushareClient
from backend.data.providers.tushare_transport import TushareError
from server import DashboardService
class FixedDatetime(datetime):
fixed_now = datetime(2026, 8, 28, 15, 4).astimezone()
@classmethod
def now(cls, tz=None):
return cls.fixed_now
class WindowClient(TushareClient):
def __init__(self, token: str = "test-token"):
super().__init__(token)
self.calls: list[str] = []
self.daily_rows: list[dict] = []
self.rt_k_error: Exception | None = None
def query(self, api_name, params=None, fields=""):
self.calls.append(api_name)
params = params or {}
if api_name == "trade_cal":
return [
{
"cal_date": "20260828",
"is_open": 1,
"pretrade_date": "20260827",
}
]
if api_name == "daily":
return list(self.daily_rows)
if api_name == "rt_k":
if self.rt_k_error is not None:
raise self.rt_k_error
raise AssertionError("rt_k should not be called in this scenario")
if api_name in {"limit_list_d", "stock_basic", "stk_limit", "daily_basic"}:
return []
raise AssertionError(f"Unexpected API call: {api_name} {params}")
def resolve_trade_context(self, requested_date: str):
return requested_date, "20260827"
class SyncDatabaseStub:
def __init__(self, latest=None):
self.latest = latest
self.snapshots: dict[str, dict] = {}
self.sync_runs: list[dict] = []
self._sync_id = 0
def start_sync(self, trade_date: str, source: str) -> int:
self._sync_id += 1
self.sync_runs.append(
{
"id": self._sync_id,
"trade_date": trade_date,
"source": source,
"status": "running",
}
)
return self._sync_id
def finish_sync(
self,
sync_id: int,
status: str,
record_count: int = 0,
message: str = "",
source: str | None = None,
) -> None:
for row in self.sync_runs:
if row["id"] == sync_id:
row.update(
{
"status": status,
"record_count": record_count,
"message": message,
"source": source or row["source"],
}
)
return
raise AssertionError(f"unknown sync_id {sync_id}")
def save_snapshot(self, trade_date: str, source: str, payload: dict) -> None:
self.snapshots[trade_date] = {"source": source, "payload": payload}
def save_data_snapshot(self, kind: str, cache_key: str, source: str, payload: dict) -> None:
return None
def get_latest_real_snapshot(self, _trade_date: str, strictly_before: bool = False):
return self.latest
def reason_overrides(self, _trade_date: str):
return {}
class DashboardRefreshWindowTests(unittest.TestCase):
def setUp(self) -> None:
TushareClient._realtime_reference_cache.clear()
TushareClient._capital_cache.clear()
TushareClient._latest_realtime_market.clear()
TushareClient._stock_activity_cache.clear()
def _service(self, client: WindowClient, latest=None) -> DashboardService:
service = object.__new__(DashboardService)
service._system_credentials = {"tushare_token": "test-token"}
service.sync_lock = threading.Lock()
service.database = SyncDatabaseStub(latest=latest)
service.data_gateway = None
service._tushare_client = lambda: client
service._enrich_dashboard_sentiment = lambda dashboard, _date: dashboard
service._apply_reason_overrides = lambda dashboard: dashboard
return service
def test_should_use_realtime_at_1504(self) -> None:
FixedDatetime.fixed_now = datetime(2026, 8, 28, 15, 4).astimezone()
with patch("backend.data.providers.tushare_dashboard.datetime", FixedDatetime):
self.assertTrue(TushareClient.should_use_realtime("20260828", "20260828"))
def test_should_use_daily_at_1505(self) -> None:
FixedDatetime.fixed_now = datetime(2026, 8, 28, 15, 5).astimezone()
with patch("backend.data.providers.tushare_dashboard.datetime", FixedDatetime):
self.assertFalse(TushareClient.should_use_realtime("20260828", "20260828"))
def test_after_close_empty_daily_does_not_call_rt_k(self) -> None:
FixedDatetime.fixed_now = datetime(2026, 8, 28, 15, 49).astimezone()
client = WindowClient()
client.daily_rows = []
with patch("backend.data.providers.tushare_dashboard.datetime", FixedDatetime):
with self.assertRaises(TushareError):
client.dashboard("20260828")
self.assertIn("daily", client.calls)
self.assertNotIn("rt_k", client.calls)
def test_after_close_uses_daily_when_ready(self) -> None:
FixedDatetime.fixed_now = datetime(2026, 8, 28, 15, 49).astimezone()
client = WindowClient()
client.daily_rows = [
{
"ts_code": "000001.SZ",
"trade_date": "20260828",
"open": 10,
"high": 11,
"low": 9.5,
"close": 10.5,
"pre_close": 10,
"pct_chg": 5,
"vol": 1000,
"amount": 1_000_000,
}
]
def load_limit_lists(_trade_date):
return []
def load_limit_type(_trade_date, _limit_type):
return []
client._load_limit_lists = load_limit_lists # type: ignore[method-assign]
client._load_limit_type = load_limit_type # type: ignore[method-assign]
client._derive_limits = lambda *args, **kwargs: [] # type: ignore[method-assign]
with patch("backend.data.providers.tushare_dashboard.datetime", FixedDatetime):
dashboard = client.dashboard("20260828")
self.assertIn("daily", client.calls)
self.assertNotIn("rt_k", client.calls)
self.assertFalse(dashboard["meta"].get("realtime"))
self.assertEqual(dashboard["meta"]["trade_date"], "2026-08-28")
def test_fallback_old_snapshot_marks_sync_and_job_status_failed(self) -> None:
FixedDatetime.fixed_now = datetime(2026, 8, 28, 15, 49).astimezone()
client = WindowClient()
client.daily_rows = []
latest = {
"meta": {"source": "tushare", "trade_date": "2026-08-27"},
"overview": {},
"limits": [],
"broken": [],
"down_limits": [],
"yesterday_limits": [],
}
service = self._service(client, latest=latest)
with patch("backend.data.providers.tushare_dashboard.datetime", FixedDatetime):
result = service.sync_dashboard("20260828")
self.assertEqual(result["status"], "failed")
self.assertTrue(result["meta"]["carried_forward"])
self.assertEqual(result["meta"]["trade_date"], "2026-08-27")
self.assertEqual(service.database.sync_runs[-1]["status"], "failed")
self.assertNotIn("rt_k", client.calls)
def test_closing_snapshot_due_after_1505_when_today_missing(self) -> None:
FixedDatetime.fixed_now = datetime(2026, 8, 28, 15, 49).astimezone()
service = object.__new__(DashboardService)
service._system_credentials = {"tushare_token": "test-token"}
with patch("backend.features.market.service.datetime", FixedDatetime), patch(
"backend.features.market.service.date"
) as fake_date:
fake_date.today.return_value = FixedDatetime.fixed_now.date()
self.assertTrue(service._closing_snapshot_due("20260828", {}))
self.assertFalse(
service._closing_snapshot_due(
"20260828",
{
"meta": {
"trade_date": "2026-08-28",
"realtime": False,
}
},
)
)
if __name__ == "__main__":
unittest.main()
+317
View File
@@ -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()
+3
View File
@@ -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 <admin> [--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;
+112
View File
@@ -0,0 +1,112 @@
#!/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
import sys
from datetime import date
from pathlib import Path
ROOT = Path(__file__).resolve().parents[1]
if str(ROOT) not in sys.path:
sys.path.insert(0, str(ROOT))
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()