From 0a5a14cbbb5080b1ffd53d67d5101e307df56b9d Mon Sep 17 00:00:00 2001 From: shangfangjian Date: Mon, 21 Sep 2026 01:36:54 +0800 Subject: [PATCH] =?UTF-8?q?=E9=87=87=E9=9B=86=E4=BB=BB=E5=8A=A1=E6=94=AF?= =?UTF-8?q?=E6=8C=81=E6=89=8B=E5=8A=A8=E8=A7=A6=E5=8F=91=20+=20=E5=AE=9A?= =?UTF-8?q?=E6=97=B6=E8=B0=83=E5=BA=A6(cron)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 后端: - 新增 Scheduler + ScheduledTask 核心模块(croniter 解析 cron 表达式) - 新增 Schedule 数据模型(alembic 0016) - 新增 /api/v1/admin/schedules CRUD 接口 - lifespan 启动时注册默认定时任务并加载数据库配置 - 默认任务:每日比赛/积分榜/30分钟统计回填(默认禁用) 前端: - 设置页新增「定时任务」管理区块 - 支持:新建/启用禁用/立即执行/删除定时任务 - 显示:任务ID、类型、cron表达式、上次执行时间/状态 依赖:新增 croniter>=2.0 Co-Authored-By: new-provider/LongCat-2.0 <> --- alembic/versions/0016_schedules.py | 38 ++++++++ frontend/src/admin/dal.ts | 32 ++++++ frontend/src/admin/pages/Settings.tsx | 77 +++++++++++++++ pyproject.toml | 1 + src/api/app.py | 34 +++++++ src/api/routes/schedules.py | 134 ++++++++++++++++++++++++++ src/api/schemas.py | 18 ++++ src/core/scheduler.py | 112 +++++++++++++++++++++ src/db/models.py | 15 +++ 9 files changed, 461 insertions(+) create mode 100644 alembic/versions/0016_schedules.py create mode 100644 src/api/routes/schedules.py create mode 100644 src/core/scheduler.py diff --git a/alembic/versions/0016_schedules.py b/alembic/versions/0016_schedules.py new file mode 100644 index 0000000..f4a330a --- /dev/null +++ b/alembic/versions/0016_schedules.py @@ -0,0 +1,38 @@ +"""新增 schedules 表 + +Revision ID: 0016_schedules +Revises: 0015_bzzoiro_single_source +Create Date: 2026-09-21 + +定时采集任务配置表。 +""" + +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + +# revision identifiers, used by Alembic. +revision: str = '0016_schedules' +down_revision: Union[str, None] = '0015_bzzoiro_single_source' +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.create_table( + 'schedules', + sa.Column('id', sa.String(50), primary_key=True), + sa.Column('task', sa.String(20), nullable=False), + sa.Column('cron', sa.String(100), nullable=False), + sa.Column('leagues', sa.Text()), + sa.Column('enabled', sa.Boolean(), server_default='true'), + sa.Column('last_run_at', sa.DateTime(timezone=True)), + sa.Column('last_status', sa.String(20)), + sa.Column('created_at', sa.DateTime(timezone=True)), + sa.Column('updated_at', sa.DateTime(timezone=True)), + ) + + +def downgrade() -> None: + op.drop_table('schedules') diff --git a/frontend/src/admin/dal.ts b/frontend/src/admin/dal.ts index 13e2bc5..0f496da 100644 --- a/frontend/src/admin/dal.ts +++ b/frontend/src/admin/dal.ts @@ -406,3 +406,35 @@ export async function fetchKeyRingStatus(): Promise { export async function resetKeyRingCooldown(): Promise<{ ok: boolean; message: string; stats: KeyRingStatusResponse }> { return api.post<{ ok: boolean; message: string; stats: KeyRingStatusResponse }>(`${API_BASE}/admin/keyring/cooldown/reset`) } + +// ── 定时任务 ──────────────────────────────────────────────────── + +export interface ScheduleItem { + id: string + task: string + cron: string + leagues?: string + enabled: boolean + last_run_at?: string | null + last_status?: string | null +} + +export async function fetchSchedules(): Promise { + return api.get(`${API_BASE}/admin/schedules`) +} + +export async function createSchedule(data: { id: string; task: string; cron: string; leagues?: string; enabled: boolean }): Promise<{ ok: boolean }> { + return api.post<{ ok: boolean }>(`${API_BASE}/admin/schedules`, data) +} + +export async function updateSchedule(id: string, data: Partial): Promise<{ ok: boolean }> { + return api.put<{ ok: boolean }>(`${API_BASE}/admin/schedules/${id}`, data) +} + +export async function deleteSchedule(id: string): Promise<{ ok: boolean }> { + return api.delete<{ ok: boolean }>(`${API_BASE}/admin/schedules/${id}`) +} + +export async function runScheduleNow(id: string): Promise<{ ok: boolean; message: string }> { + return api.post<{ ok: boolean; message: string }>(`${API_BASE}/admin/schedules/${id}/run`) +} diff --git a/frontend/src/admin/pages/Settings.tsx b/frontend/src/admin/pages/Settings.tsx index e0d8cff..e22ff62 100644 --- a/frontend/src/admin/pages/Settings.tsx +++ b/frontend/src/admin/pages/Settings.tsx @@ -13,7 +13,9 @@ import { fetchSettings, updateSetting, clearSetting, testLLMConnection, fetchLLMUsageStats, fetchLLMModels, fetchKeyRingStatus, resetKeyRingCooldown, + fetchSchedules, createSchedule, updateSchedule, deleteSchedule, runScheduleNow, } from '../dal' +import type { ScheduleItem } from '../dal' import { changePassword, fetchAuthState, UNAUTHORIZED_EVENT } from '../api' import type { LLMUsageStats, DataSourceSetting } from '../types' import type { KeyRingStatusResponse } from '../dal' @@ -49,6 +51,10 @@ export default function SettingsPage() { const [pwdBusy, setPwdBusy] = useState(false) const [pwdNotice, setPwdNotice] = useState<{ ok: boolean; text: string } | null>(null) + // 定时任务 + const [schedules, setSchedules] = useState([]) + const [schedulesLoading, setSchedulesLoading] = useState(true) + // ── 数据加载 ── const loadSettings = useCallback(async () => { setSettingsLoading(true) @@ -89,6 +95,7 @@ export default function SettingsPage() { loadLlmStats() loadKeyRing() fetchAuthState().then(s => setPasswordOrigin(s.password_origin ?? null)).catch(() => {}) + fetchSchedules().then(setSchedules).catch(() => []).finally(() => setSchedulesLoading(false)) }, [loadSettings, loadLlmStats, loadKeyRing]) const dataSourceSettings = allSettings.filter(s => DATA_SOURCE_KEYS.includes(s.key)) @@ -362,6 +369,76 @@ export default function SettingsPage() { + + {/* ── 4. 定时任务 ── */} +
+

定时任务

+ + { + createSchedule({ id: `schedule-${Date.now()}`, task: 'events', cron: '0 8 * * *', leagues: undefined, enabled: false }) + .then(() => fetchSchedules().then(setSchedules)) + }} + className="btn btn-sm" + > + + 新建 + + } + /> + + {schedulesLoading ? ( + + ) : schedules.length === 0 ? ( +

暂无定时任务,点击右上角「+ 新建」创建

+ ) : ( +
+ {schedules.map(s => ( +
+
+
+ + {s.id} + {s.task} +
+

{s.cron}

+ {s.last_run_at && ( +

+ 上次: {new Date(s.last_run_at).toLocaleString('zh-CN', { hour12: false })} + {s.last_status === 'success' ? ' ✓' : s.last_status === 'failed' ? ' ✗' : ''} +

+ )} +
+
+ + + +
+
+ ))} +
+ )} +
+
+
) } diff --git a/pyproject.toml b/pyproject.toml index 3f39cd3..54aa03d 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -14,6 +14,7 @@ dependencies = [ "httpx>=0.27", "alembic>=1.13", "cryptography>=42.0", + "croniter>=2.0", ] [project.optional-dependencies] diff --git a/src/api/app.py b/src/api/app.py index 06c7643..cebaa41 100644 --- a/src/api/app.py +++ b/src/api/app.py @@ -22,11 +22,43 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]: migrate_plaintext_sensitive_settings, ) from src.core.security_check import assert_security_on_startup + from src.core.scheduler import scheduler + from src.api.routes.schedules import _run_scheduled_task + from src.data.config import BZZOIRO_LEAGUE_IDS await init_db() # 验证连接,不建表 await migrate_plaintext_sensitive_settings() # 明文敏感配置 → 加密(幂等) await ensure_admin_password_hashed() # .env 明文密码 → scrypt 哈希(幂等) await assert_security_on_startup() # 启动安全校验(生产拒绝/开发警告) + + # 注册默认定时任务(如果数据库中没有) + from src.db.base import AsyncSessionLocal + from sqlalchemy import select + from src.db.models import Schedule + async with AsyncSessionLocal() as session: + existing = (await session.execute(select(Schedule.id))).scalars().all() + if "daily-events" not in existing: + session.add(Schedule(id="daily-events", task="events", cron="0 8 * * *", enabled=False)) + if "daily-standings" not in existing: + session.add(Schedule(id="daily-standings", task="standings", cron="0 9 * * *", enabled=False)) + if "daily-stats" not in existing: + session.add(Schedule(id="daily-stats", task="stats", cron="*/30 * * * *", enabled=False)) + await session.commit() + + # 从数据库加载所有启用的定时任务 + async with AsyncSessionLocal() as session: + schedules = (await session.execute(select(Schedule).where(Schedule.enabled))).scalars().all() + for s in schedules: + leagues = s.leagues.split(",") if s.leagues else list(BZZOIRO_LEAGUE_IDS.keys()) + scheduler.register( + s.id, s.task, + lambda sid=s.id: _run_scheduled_task(sid), + enabled=s.enabled, + ) + + await scheduler.start() + logger.info("应用启动完成") yield + await scheduler.stop() await close_client() @@ -62,6 +94,7 @@ def create_app() -> FastAPI: from src.api.routes.backtest import router as backtest_router from src.api.routes.auth import router as auth_router from src.api.routes.admin_settings import router as admin_settings_router + from src.api.routes.schedules import router as schedules_router app.include_router(matches_router) app.include_router(predict_router) @@ -70,6 +103,7 @@ def create_app() -> FastAPI: app.include_router(backtest_router) app.include_router(auth_router) app.include_router(admin_settings_router) + app.include_router(schedules_router) @app.get("/health") async def health(): diff --git a/src/api/routes/schedules.py b/src/api/routes/schedules.py new file mode 100644 index 0000000..c006b09 --- /dev/null +++ b/src/api/routes/schedules.py @@ -0,0 +1,134 @@ +"""定时任务管理路由。""" +from __future__ import annotations + +import logging +from datetime import datetime, timezone + +from fastapi import APIRouter, Depends, HTTPException +from sqlalchemy import select, delete + +from src.api.deps import require_admin +from src.api.schemas import ScheduleIn, ScheduleOut +from src.core.scheduler import scheduler +from src.data.bzzoiro import ingest_bzzoiro_event_stats, ingest_bzzoiro_standings +from src.data.sources import get_source +from src.data.config import BZZOIRO_LEAGUE_IDS +from src.db.base import AsyncSession, get_db_read +from src.db.models import Schedule +from src.db.unit_of_work import get_uow + +logger = logging.getLogger(__name__) + +router = APIRouter(prefix="/api/v1/admin", tags=["schedule"], dependencies=[Depends(require_admin)]) + + +async def _run_scheduled_task(schedule_id: str) -> None: + """执行定时任务的回调函数。""" + async with get_uow() as session: + stmt = select(Schedule).where(Schedule.id == schedule_id) + sched = (await session.execute(stmt)).scalar_one_or_none() + if sched is None or not sched.enabled: + return + + leagues = sched.leagues.split(",") if sched.leagues else list(BZZOIRO_LEAGUE_IDS.keys()) + task = sched.task + + try: + if task in ("events", "all"): + statuses = ["finished", "scheduled"] + source = get_source("bzzoiro") + for st in statuses: + await source.ingest(session, leagues=leagues, status=st) + + if task in ("standings", "all"): + await ingest_bzzoiro_standings(session, leagues=leagues) + + if task in ("stats", "all"): + await ingest_bzzoiro_event_stats(session, leagues=leagues, limit=500, only_missing=True) + + sched.last_status = "success" + except Exception: + logger.exception("定时任务执行失败: %s", schedule_id) + sched.last_status = "failed" + finally: + sched.last_run_at = datetime.now(timezone.utc) + + +def _sync_scheduler() -> None: + """同步数据库中的调度配置到调度器。""" + # 这是一个简化版本:实际应该在 lifespan 中异步同步 + pass + + +@router.get("/schedules") +async def list_schedules(db: AsyncSession = Depends(get_db_read)): + """列出所有定时任务。""" + rows = (await db.execute(select(Schedule).order_by(Schedule.created_at))).scalars().all() + return [ + ScheduleOut( + id=s.id, + task=s.task, + cron=s.cron, + leagues=s.leagues, + enabled=s.enabled, + last_run_at=s.last_run_at.isoformat() if s.last_run_at else None, + last_status=s.last_status, + ) + for s in rows + ] + + +@router.post("/schedules") +async def create_schedule(req: ScheduleIn, db: AsyncSession = Depends(get_db_read)): + """创建定时任务。""" + sched = Schedule( + id=req.id, + task=req.task, + cron=req.cron, + leagues=",".join(req.leagues) if req.leagues else None, + enabled=req.enabled, + ) + db.add(sched) + await db.commit() + + # 注册到调度器 + scheduler.register(req.id, req.task, lambda: _run_scheduled_task(req.id), enabled=req.enabled) + + return {"ok": True, "id": req.id} + + +@router.put("/schedules/{schedule_id}") +async def update_schedule(schedule_id: str, req: ScheduleIn, db: AsyncSession = Depends(get_db_read)): + """更新定时任务。""" + stmt = select(Schedule).where(Schedule.id == schedule_id) + sched = (await db.execute(stmt)).scalar_one_or_none() + if sched is None: + raise HTTPException(404, "定时任务不存在") + + sched.task = req.task + sched.cron = req.cron + sched.leagues = ",".join(req.leagues) if req.leagues else None + sched.enabled = req.enabled + await db.commit() + + # 更新调度器 + scheduler.register(schedule_id, req.task, lambda: _run_scheduled_task(schedule_id), enabled=req.enabled) + + return {"ok": True} + + +@router.delete("/schedules/{schedule_id}") +async def delete_schedule(schedule_id: str, db: AsyncSession = Depends(get_db_read)): + """删除定时任务。""" + await db.execute(delete(Schedule).where(Schedule.id == schedule_id)) + await db.commit() + scheduler.remove(schedule_id) + return {"ok": True} + + +@router.post("/schedules/{schedule_id}/run") +async def run_schedule_now(schedule_id: str): + """手动触发定时任务。""" + import asyncio + asyncio.create_task(_run_scheduled_task(schedule_id)) + return {"ok": True, "message": "任务已启动"} diff --git a/src/api/schemas.py b/src/api/schemas.py index ef77060..0a581bd 100644 --- a/src/api/schemas.py +++ b/src/api/schemas.py @@ -121,6 +121,24 @@ class IngestResponse(BaseModel): errors: list[str] = [] +class ScheduleIn(BaseModel): + id: str = Field(..., description="任务唯一标识,如 'daily-events'") + task: str = Field(..., description="events / standings / stats / all") + cron: str = Field(..., description="cron 表达式,如 '0 8 * * *' (每天 8 点)") + leagues: list[str] = Field(default_factory=list, description="联赛代码列表,空=全部") + enabled: bool = True + + +class ScheduleOut(BaseModel): + id: str + task: str + cron: str + leagues: str | None + enabled: bool + last_run_at: str | None + last_status: str | None + + class SettleRequest(BaseModel): prediction_id: int home_goals: int = Field(ge=0, le=30) diff --git a/src/core/scheduler.py b/src/core/scheduler.py new file mode 100644 index 0000000..a9ffc85 --- /dev/null +++ b/src/core/scheduler.py @@ -0,0 +1,112 @@ +"""定时任务调度器:支持 cron 表达式触发采集任务。""" +from __future__ import annotations + +import asyncio +import logging +from datetime import datetime, timedelta +from typing import Callable, Coroutine + +from croniter import croniter + +logger = logging.getLogger(__name__) + + +class ScheduledTask: + """一个定时任务。""" + + def __init__( + self, + task_id: str, + cron: str, + fn: Callable[[], Coroutine], + enabled: bool = True, + ) -> None: + self.task_id = task_id + self.cron = cron + self.fn = fn + self.enabled = enabled + self.last_run: datetime | None = None + self.next_run: datetime | None = None + self._task: asyncio.Task | None = None + self._calc_next() + + def _calc_next(self) -> None: + try: + self.next_run = croniter(self.cron, datetime.now).get_next(datetime) + except Exception: + self.next_run = None + + def update(self, cron: str | None = None, enabled: bool | None = None) -> None: + if cron is not None: + self.cron = cron + if enabled is not None: + self.enabled = enabled + self._calc_next() + + async def _run_loop(self) -> None: + while True: + if not self.enabled or not self.next_run: + await asyncio.sleep(60) + self._calc_next() + continue + + now = datetime.now() + wait_seconds = (self.next_run - now).total_seconds() + if wait_seconds > 0: + await asyncio.sleep(min(wait_seconds, 60)) + continue + + # 执行任务 + self.last_run = datetime.now() + self._calc_next() + try: + logger.info("定时任务触发: %s (cron=%s)", self.task_id, self.cron) + await self.fn() + logger.info("定时任务完成: %s", self.task_id) + except Exception: + logger.exception("定时任务失败: %s", self.task_id) + + +class Scheduler: + """全局定时任务调度器。""" + + def __init__(self) -> None: + self._tasks: dict[str, ScheduledTask] = {} + + def register( + self, + task_id: str, + cron: str, + fn: Callable[[], Coroutine], + enabled: bool = True, + ) -> ScheduledTask: + if task_id in self._tasks: + self._tasks[task_id].update(cron=cron, enabled=enabled) + return self._tasks[task_id] + task = ScheduledTask(task_id, cron, fn, enabled) + self._tasks[task_id] = task + return task + + def get(self, task_id: str) -> ScheduledTask | None: + return self._tasks.get(task_id) + + def list_all(self) -> list[ScheduledTask]: + return list(self._tasks.values()) + + def remove(self, task_id: str) -> None: + self._tasks.pop(task_id, None) + + async def start(self) -> None: + for task in self._tasks.values(): + task._task = asyncio.create_task(task._run_loop()) + logger.info("定时调度器已启动, 共 %d 个任务", len(self._tasks)) + + async def stop(self) -> None: + for task in self._tasks.values(): + if task._task: + task._task.cancel() + self._tasks.clear() + + +# 全局单例 +scheduler = Scheduler() diff --git a/src/db/models.py b/src/db/models.py index 7822aca..5feb4b1 100644 --- a/src/db/models.py +++ b/src/db/models.py @@ -264,6 +264,21 @@ class AppSetting(Base): updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_utcnow, onupdate=_utcnow) +class Schedule(Base): + """定时采集任务配置。""" + __tablename__ = "schedules" + + id: Mapped[str] = mapped_column(String(50), primary_key=True) + task: Mapped[str] = mapped_column(String(20), nullable=False) # events / standings / stats / all + cron: Mapped[str] = mapped_column(String(100), nullable=False) # cron 表达式 + leagues: Mapped[str | None] = mapped_column(Text) # 逗号分隔的联赛代码,空=全部 + enabled: Mapped[bool] = mapped_column(Boolean, default=True) + last_run_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True)) + last_status: Mapped[str | None] = mapped_column(String(20)) # success / failed + created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_utcnow) + updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_utcnow, onupdate=_utcnow) + + # ── 数据管线基础设施(对应迁移 0008) ────────────────────────────────── # Bronze 层、死信、质量监控、血缘追踪 4 张表。