diff --git a/app/core/db/models.py b/app/core/db/models.py index 2e4b76f..52bb4c1 100644 --- a/app/core/db/models.py +++ b/app/core/db/models.py @@ -95,6 +95,41 @@ class DatasetRegistry(ORMBase): updated_at: Mapped[Optional[datetime]] = _updated_at() +# ──────────────────────── 2b. sync_history ──────────────────────── +class SyncHistory(ORMBase): + """同步任务执行历史(每次跑都留一行,append-only)。 + + 与 dataset_registry 的区别:dataset_registry 只保留"最后一次状态",便于快查; + sync_history 保留完整历史,支持看板按天/按 task 维度聚合 + 失败回溯。 + + append-only:每次 task.run() 完成(无论 ok/warning/failed/blocked)插入一行。 + """ + __tablename__ = "sync_history" + __table_args__ = ( + Index("idx_sync_history_dataset_started", "dataset_id", "started_at"), + Index("idx_sync_history_started", "started_at"), + Index("idx_sync_history_status", "status"), + {"schema": "market_data"}, + ) + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + dataset_id: Mapped[str] = mapped_column(String(64), nullable=False) + run_date: Mapped[date] = mapped_column(Date, nullable=False) + status: Mapped[str] = mapped_column(String(16), nullable=False) # ok/warning/error/blocked + trigger_source: Mapped[str] = mapped_column(String(32), nullable=False, default="manual") + started_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False) + finished_at: Mapped[Optional[datetime]] = mapped_column(DateTime(timezone=True)) + elapsed_sec: Mapped[Optional[float]] = mapped_column(Float, default=0) + rows_written: Mapped[int] = mapped_column(Integer, default=0) + message: Mapped[Optional[str]] = mapped_column(Text) + error: Mapped[Optional[str]] = mapped_column(Text) + # 子任务维度统计(JSONB): {"ok": 5204, "fail": 1, "skip": 0, "total": 5205} + stats: Mapped[Optional[dict]] = mapped_column(JSONB, default=dict) + # 触发者上下文 (cli 命令 / scheduler key / manual user) + triggered_by: Mapped[Optional[str]] = mapped_column(String(64), default="") + updated_at: Mapped[Optional[datetime]] = _updated_at() + + # ──────────────────────── 3. stocks ──────────────────────── class Stock(ORMBase): """股票基础信息 + 最新股本快照(冗余缓存)。""" diff --git a/app/core/db/ops.py b/app/core/db/ops.py index 4a3c806..6862d02 100644 --- a/app/core/db/ops.py +++ b/app/core/db/ops.py @@ -12,7 +12,8 @@ from __future__ import annotations import json -from datetime import date, datetime +import os +from datetime import date, datetime, timedelta from typing import Any, Iterable, Optional from sqlalchemy import and_, case, delete, exists, func, or_, select, update @@ -40,6 +41,7 @@ from app.core.db.models import ( Stock, StockNodeMap, StockSectorMap, + SyncHistory, TickTrade, ) from app.core.db.orm import get_session @@ -98,10 +100,31 @@ def _row_to_dict(row) -> dict[str, Any]: def _ensure_schema() -> None: """所有读函数开头调用一次,幂等。 - PG 的表由 ``app.core.db.pg_bootstrap`` 一次性建好(人工触发),这里 no-op 即可。 - 保留这个函数是为了不破坏 9 个 sync task 的调用约定(行 14 个调用点)。 + 用 superuser 连接跑 create_all(market_sync role 没 CREATE 权限)。 + SA 2.x 内置 IF NOT EXISTS 检查,已存在则跳过。 + 这样新增 model 后无需手动跑 pg_bootstrap。 + + 首次部署仍然建议跑 pg_bootstrap(它还要建 role / schema); + 这里是后续迭代加表的兜底。 """ - pass + from app.core.db import models # noqa: F401 触发全部 model 注册 + from app.core.config import settings + from app.core.db.orm import ORMBase + from sqlalchemy import create_engine + + su_url = ( + f"postgresql+psycopg2://{os.environ.get('PG_SUPERUSER', 'postgres')}:" + f"{os.environ.get('PG_SUPERUSER_PASSWORD', 'postgres')}@" + f"{os.environ.get('PG_SUPERHOST', settings.pg_host)}:" + f"{os.environ.get('PG_SUPERPORT', settings.pg_port)}/{settings.pg_db_name}" + ) + su_engine = create_engine(su_url, future=True) + try: + ORMBase.metadata.create_all(su_engine) + except Exception as e: + # 创建表失败不应阻塞主流程(表可能已存在 / 权限不够) + import logging + logging.getLogger("sync").debug(f"[_ensure_schema] create_all 失败(忽略): {e}") # ── config 表 ─────────────────────────────────────────────────────────── @@ -281,6 +304,123 @@ def update_dataset_registry_state(dataset_id: str, **kwargs: Any) -> None: ) +# ──────── sync_history (append-only 每次同步留痕) ──────── + + +def insert_sync_history( + *, + dataset_id: str, + run_date: date, + status: str, + trigger_source: str, + started_at: datetime, + finished_at: Optional[datetime] = None, + elapsed_sec: float = 0.0, + rows_written: int = 0, + message: str = "", + error: str = "", + stats: Optional[dict] = None, + triggered_by: str = "", +) -> int: + """每次 task.run() 完成追加一行 sync_history。 + + 看板和 MCP 都从这张表读历史;失败排查也靠它。 + 返回新行的 id。 + """ + row = SyncHistory( + dataset_id=dataset_id, + run_date=run_date, + status=status, + trigger_source=trigger_source, + started_at=started_at, + finished_at=finished_at, + elapsed_sec=float(elapsed_sec), + rows_written=int(rows_written), + message=message or None, + error=error or None, + stats=stats or {}, + triggered_by=triggered_by or None, + ) + with get_session() as s: + s.add(row) + s.flush() + return int(row.id) + + +def list_sync_history( + *, + dataset_id: Optional[str] = None, + run_date: Optional[date] = None, + days: int = 7, + limit: int = 200, +) -> list[dict[str, Any]]: + """查 sync_history(看板 / API 用),按 started_at DESC。 + + 默认最近 7 天。指定 dataset_id 时只查该任务。 + """ + with get_session() as s: + q = select(SyncHistory) + if dataset_id: + q = q.where(SyncHistory.dataset_id == dataset_id) + if run_date: + q = q.where(SyncHistory.run_date == run_date) + else: + cutoff = date.today() - timedelta(days=days) + q = q.where(SyncHistory.run_date >= cutoff) + q = q.order_by(SyncHistory.started_at.desc()).limit(limit) + rows = s.execute(q).scalars().all() + return [ + { + "id": r.id, + "dataset_id": r.dataset_id, + "run_date": r.run_date.isoformat() if r.run_date else None, + "status": r.status, + "trigger_source": r.trigger_source, + "started_at": r.started_at.isoformat() if r.started_at else None, + "finished_at": r.finished_at.isoformat() if r.finished_at else None, + "elapsed_sec": r.elapsed_sec, + "rows_written": r.rows_written, + "message": r.message, + "error": r.error, + "stats": r.stats, + "triggered_by": r.triggered_by, + } + for r in rows + ] + + +def daily_sync_summary(days: int = 7) -> list[dict[str, Any]]: + """按日期聚合(看板首页用):每天 ok/warning/error/blocked 计数。 + + 返回 [{date, ok, warning, error, blocked, total, tasks_run}, ...] + """ + with get_session() as s: + cutoff = date.today() - timedelta(days=days) + rows = s.execute( + select( + SyncHistory.run_date, + SyncHistory.status, + func.count().label("n"), + func.count(func.distinct(SyncHistory.dataset_id)).label("tasks"), + ) + .where(SyncHistory.run_date >= cutoff) + .group_by(SyncHistory.run_date, SyncHistory.status) + .order_by(SyncHistory.run_date.desc()) + ).all() + # pivot 成每天一行 + by_date: dict[date, dict[str, Any]] = {} + for run_date, status, n, tasks in rows: + d = by_date.setdefault(run_date, { + "date": run_date.isoformat(), + "ok": 0, "warning": 0, "error": 0, "blocked": 0, + "total": 0, "tasks_run": 0, + }) + d[status] = int(n) + d["total"] += int(n) + d["tasks_run"] = max(d["tasks_run"], int(tasks)) # max 而不是 sum(去重) + return sorted(by_date.values(), key=lambda x: x["date"], reverse=True) + + def recover_interrupted_dataset_registry() -> int: """把状态卡在 running 的同步任务标记为 failed(启动时调用)。""" with get_session() as s: diff --git a/app/core/sync/base.py b/app/core/sync/base.py index 87c535d..fc5280d 100644 --- a/app/core/sync/base.py +++ b/app/core/sync/base.py @@ -4,6 +4,7 @@ - 状态自动更新(running → success / failed) - 进度回调 - 异常捕获 → 写 dataset_registry +- 每次 run 完成追加一行 sync_history(看板/MCP/失败排查) 子类只需实现 _run() 即可。 """ @@ -14,6 +15,7 @@ import traceback from datetime import datetime, time as dtime from typing import Any, Optional +from app.core.db import ops as db_ops from app.core.sync.registry import ( mark_sync_failed, mark_sync_progress, @@ -55,7 +57,10 @@ class SyncTask: # ── 公开入口 ───────────────────────────────────────── def run(self, *, trigger_source: str = "manual", **kwargs) -> dict[str, Any]: - """统一入口:自动包 mark_running / mark_success / mark_failed。""" + """统一入口:自动包 mark_running / mark_success / mark_failed。 + + 完成后追加一行 sync_history(看板 / MCP / 失败排查共用)。 + """ # 1. 确保 health check 跑过(CLI / 手动触发场景下 schedule 后台不会先跑) self._ensure_health_checked() @@ -66,6 +71,7 @@ class SyncTask: message=f"开始同步 {self.dataset_id}...", ) t0 = time.time() + started_at = datetime.now() try: result = self._run(trigger_source=trigger_source, **kwargs) elapsed = round(time.time() - t0, 1) @@ -94,6 +100,17 @@ class SyncTask: logger.info(f"[{self.dataset_id}] 完成 {status}: {message} ({elapsed}s)") # 同步 schedule config 的 lastRun(fix Bug 2) self._maybe_mirror_schedule_status(trigger_source, status, message) + # 追加 sync_history(每次 run 留痕,看板/MCP/失败排查共用) + self._record_history( + trigger_source=trigger_source, + started_at=started_at, + finished_at=datetime.now(), + elapsed_sec=elapsed, + status=status, + message=message, + error=result.get("error", ""), + result=result, + ) return result except Exception as e: elapsed = round(time.time() - t0, 1) @@ -101,6 +118,17 @@ class SyncTask: err_msg = f"{type(e).__name__}: {e}" mark_sync_failed(self.dataset_id, message=err_msg, error=tb) logger.error(f"[{self.dataset_id}] 失败: {err_msg}\n{tb}") + # 异常路径也要写 sync_history + self._record_history( + trigger_source=trigger_source, + started_at=started_at, + finished_at=datetime.now(), + elapsed_sec=elapsed, + status="error", + message=err_msg, + error=tb, + result={"status": "error"}, + ) return { "status": "error", "message": err_msg, @@ -108,6 +136,45 @@ class SyncTask: "elapsed_sec": elapsed, } + def _record_history( + self, *, + trigger_source: str, + started_at: datetime, + finished_at: datetime, + elapsed_sec: float, + status: str, + message: str, + error: str, + result: dict[str, Any], + ) -> None: + """把这次 run 的核心指标落 sync_history(看板/MCP 共用)。 + + stats 字段收集 _run() 返回 dict 里的统计字段(ok/fail/skip/total/rows 等)。 + """ + try: + stats_keys = {"ok", "fail", "skip", "total", "rows", "rows_written", + "categories", "nodes", "mappings", "days", "sectors"} + stats = {k: result[k] for k in stats_keys if k in result} + rows_written = int(stats.get("rows", 0) or stats.get("rows_written", 0) + or stats.get("mappings", 0) or 0) + db_ops.insert_sync_history( + dataset_id=self.dataset_id, + run_date=started_at.date(), + status=status, + trigger_source=trigger_source, + started_at=started_at, + finished_at=finished_at, + elapsed_sec=elapsed_sec, + rows_written=rows_written, + message=message[:1000] if message else "", + error=error[:2000] if error else "", + stats=stats, + triggered_by=trigger_source, + ) + except Exception as e: + # 落库失败不应阻塞主流程 + logger.warning(f"[{self.dataset_id}] sync_history 落库失败(不影响主流程): {e}") + def _ensure_health_checked(self) -> None: """若内存里 health_results 为空,跑一次。""" try: diff --git a/tests/test_schema_models.py b/tests/test_schema_models.py index 96d2826..394941b 100644 --- a/tests/test_schema_models.py +++ b/tests/test_schema_models.py @@ -7,14 +7,14 @@ from __future__ import annotations def test_orm_metadata_registers_all_tables(): - """ORMBase.metadata 应该注册 22 张 PG 表(项目主业务表)。""" + """ORMBase.metadata 应该注册 23 张 PG 表(项目主业务表)。""" from app.core.db import models # 触发全部模型 import from app.core.db.orm import ORMBase # 注意:SA 2.x 的 MetaData.tables 既能按 bare name 也能按 schema-qualified name 索引 # (keyed 容器),所以用 values() 拿 Table 对象再用 .name 拿 bare 名 bare_names = {t.name for t in ORMBase.metadata.tables.values()} - assert len(bare_names) == 22, f"期望 22 张 ORM 表,实际 {len(bare_names)}: {bare_names}" + assert len(bare_names) == 23, f"期望 23 张 ORM 表,实际 {len(bare_names)}: {bare_names}" expected = { "config", "dataset_registry", "stocks", "indices", "kline_stock", "kline_index", "kline_5min", "moneyflow", "share", @@ -24,6 +24,7 @@ def test_orm_metadata_registers_all_tables(): "longhubang_daily", "longhubang_seat", # 2026-07-01 龙虎榜 "kline_stock_ma_daily", # 2026-07-04 mairui_ma_daily "node_categories", "nodes", "stock_node_map", # 2026-07-02 股票-节点映射 + "sync_history", # 2026-07-07 同步历史表 } assert bare_names == expected, f"ORM 表名集合与期望不符: {bare_names ^ expected}"