"""定时任务管理路由。""" 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, ScheduleUpdate, ScheduleOut from src.core.scheduler import scheduler from src.data.bzzoiro_standings import ingest_bzzoiro_standings from src.data.bzzoiro_stats import ingest_bzzoiro_event_stats 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() # 注册到调度器(注意: 第 2 参是 cron 表达式,不是 task 类型名) scheduler.register(req.id, req.cron, 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: ScheduleUpdate, 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, "定时任务不存在") if req.task is not None: sched.task = req.task if req.cron is not None: sched.cron = req.cron if req.leagues is not None: sched.leagues = ",".join(req.leagues) if req.leagues else None if req.enabled is not None: sched.enabled = req.enabled await db.commit() # 更新调度器(第 2 参传 cron 表达式;使用最终值) scheduler.register(schedule_id, sched.cron, lambda: _run_scheduled_task(schedule_id), enabled=sched.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": "任务已启动"} # ── 采集失败重试 ──────────────────────────────────────────────── @router.get("/ingest-failures") async def list_ingest_failures(db: AsyncSession = Depends(get_db_read)): """列出采集失败记录。""" from src.db.models import IngestFailure rows = ( await db.execute( select(IngestFailure).order_by(IngestFailure.created_at.desc()).limit(50) ) ).scalars().all() return [ { "id": f.id, "source": f.source_system, "entity_type": f.entity_type, "source_record_id": f.source_record_id, "error_type": f.error_type, "error_detail": f.error_detail, "retry_count": f.retry_count, "status": f.status, "next_retry_at": f.next_retry_at.isoformat() if f.next_retry_at else None, "created_at": f.created_at.isoformat() if f.created_at else None, } for f in rows ] @router.post("/ingest-failures/{failure_id}/retry") async def retry_ingest_failure(failure_id: int, db: AsyncSession = Depends(get_db_read)): """重试一次采集失败。""" from src.db.models import IngestFailure stmt = select(IngestFailure).where(IngestFailure.id == failure_id) failure = (await db.execute(stmt)).scalar_one_or_none() if failure is None: raise HTTPException(404, "失败记录不存在") failure.status = "retrying" failure.retry_count += 1 await db.commit() # 触发重试(简化版:仅标记状态,实际重试逻辑由调度器处理) return {"ok": True, "message": f"已标记重试 (第 {failure.retry_count} 次)"}