diff --git a/frontend/src/admin/pages/EvalPage.tsx b/frontend/src/admin/pages/EvalPage.tsx
index 220c2bf..1d4bc01 100644
--- a/frontend/src/admin/pages/EvalPage.tsx
+++ b/frontend/src/admin/pages/EvalPage.tsx
@@ -163,7 +163,7 @@ export default function EvalPage() {
) : summary.length === 0 ? (
-
+
) : (
AsyncIterator[None]:
for s in schedules:
leagues = s.leagues.split(",") if s.leagues else list(BZZOIRO_LEAGUE_IDS.keys())
scheduler.register(
- s.id, s.task,
+ s.id, s.cron,
lambda sid=s.id: _run_scheduled_task(sid),
enabled=s.enabled,
)
diff --git a/src/api/routes/schedules.py b/src/api/routes/schedules.py
index a00327f..75f1992 100644
--- a/src/api/routes/schedules.py
+++ b/src/api/routes/schedules.py
@@ -91,8 +91,8 @@ async def create_schedule(req: ScheduleIn, db: AsyncSession = Depends(get_db_rea
db.add(sched)
await db.commit()
- # 注册到调度器
- scheduler.register(req.id, req.task, lambda: _run_scheduled_task(req.id), enabled=req.enabled)
+ # 注册到调度器(注意: 第 2 参是 cron 表达式,不是 task 类型名)
+ scheduler.register(req.id, req.cron, lambda: _run_scheduled_task(req.id), enabled=req.enabled)
return {"ok": True, "id": req.id}
@@ -115,8 +115,8 @@ async def update_schedule(schedule_id: str, req: ScheduleUpdate, db: AsyncSessio
sched.enabled = req.enabled
await db.commit()
- # 更新调度器(使用最终值)
- scheduler.register(schedule_id, sched.task, lambda: _run_scheduled_task(schedule_id), enabled=sched.enabled)
+ # 更新调度器(第 2 参传 cron 表达式;使用最终值)
+ scheduler.register(schedule_id, sched.cron, lambda: _run_scheduled_task(schedule_id), enabled=sched.enabled)
return {"ok": True}
diff --git a/src/core/scheduler.py b/src/core/scheduler.py
index ac057e5..214bc30 100644
--- a/src/core/scheduler.py
+++ b/src/core/scheduler.py
@@ -14,6 +14,9 @@ logger = logging.getLogger(__name__)
class ScheduledTask:
"""一个定时任务。"""
+ #: 运行循环的轮询粒度(秒)。同时决定运行期 cron/next_run 变更的生效延迟上限。
+ SLEEP_TICK: float = 1.0
+
def __init__(
self,
task_id: str,
@@ -28,13 +31,29 @@ class ScheduledTask:
self.last_run: datetime | None = None
self.next_run: datetime | None = None
self._task: asyncio.Task | None = None
- self._calc_next()
+ # 构造时立即校验 cron:非法表达式直接报错,不留给 _run_loop 静默吞掉。
+ # (全量审查 C1: 传入任务类型字符串时 croniter 抛异常被吞 → 任务永不触发)
+ self._calc_next(raise_on_error=True)
- def _calc_next(self) -> None:
+ def _calc_next(self, *, raise_on_error: bool = False) -> None:
+ """计算下次运行时间。
+
+ raise_on_error=True 时非法 cron 抛 ValueError(构造期用);
+ 否则仅告警并置 next_run=None(运行期容错)。
+ """
try:
- self.next_run = croniter(self.cron, datetime.now).get_next(datetime)
- except Exception:
+ # 必须调用 datetime.now():传方法对象本身会让 croniter 在
+ # (start_time or now) 的算术里抛 TypeError,导致 next_run 永远为 None。
+ self.next_run = croniter(self.cron, datetime.now()).get_next(datetime)
+ except Exception as e:
self.next_run = None
+ msg = (
+ f"任务 {self.task_id!r} 的 cron 表达式非法: {self.cron!r} ({e})。"
+ "注意: 此处应为 cron 表达式(如 '0 8 * * *'),不是任务类型名。"
+ )
+ if raise_on_error:
+ raise ValueError(msg) from e
+ logger.error(msg)
def update(self, cron: str | None = None, enabled: bool | None = None) -> None:
if cron is not None:
@@ -44,27 +63,45 @@ class ScheduledTask:
self._calc_next()
async def _run_loop(self) -> None:
+ """任务运行循环。
+
+ 睡眠策略: 使用固定的短 tick(SLEEP_TICK 秒)轮询 next_run,而不是
+ 一次性 sleep 到 next_run。原因是运行期可通过 API 更新 cron / 手动
+ 调整 next_run;若按 wait_seconds 长时间沉睡,变更最长要等
+ wait_seconds 才生效(实测可达 60s),表现为「改了不生效」。
+ """
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)
+ if not self.enabled or not self.next_run:
+ await asyncio.sleep(self.SLEEP_TICK)
+ self._calc_next()
+ continue
+
+ now = datetime.now()
+ wait_seconds = (self.next_run - now).total_seconds()
+ if wait_seconds > 0:
+ # 短 tick 轮询,保证 next_run/cron 变更能及时被感知
+ await asyncio.sleep(min(wait_seconds, self.SLEEP_TICK))
+ continue
+
+ # 执行任务:先推进 next_run 再执行,避免任务耗时导致重复触发
+ 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 asyncio.CancelledError:
+ raise
+ except Exception:
+ # 单个任务失败不能终止循环,否则一次异常即永久停摆
+ logger.exception("定时任务失败: %s", self.task_id)
+ except asyncio.CancelledError:
+ raise
except Exception:
- logger.exception("定时任务失败: %s", self.task_id)
+ # 循环体自身的意外异常(如 _calc_next)也不能终止调度
+ logger.exception("调度循环异常: %s", self.task_id)
+ await asyncio.sleep(self.SLEEP_TICK)
class DataQualityScheduler:
@@ -152,6 +189,7 @@ class Scheduler:
def __init__(self) -> None:
self._tasks: dict[str, ScheduledTask] = {}
+ self._running = False
def register(
self,
@@ -160,13 +198,31 @@ class Scheduler:
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
+ """注册(或更新)一个定时任务。
+
+ 若调度器已 start,新任务会立即启动其运行循环 ——
+ 否则运行期通过 API 新建的任务永远不会被执行(全量审查 C1 缺陷 2)。
+ """
+ existing = self._tasks.get(task_id)
+ if existing is not None:
+ existing.update(cron=cron, enabled=enabled)
+ # 更新 cron 后重新校验:改坏了要立刻报错,而不是静默失活
+ existing._calc_next(raise_on_error=True)
+ existing.fn = fn
+ task = existing
+ else:
+ task = ScheduledTask(task_id, cron, fn, enabled)
+ self._tasks[task_id] = task
+
+ if self._running:
+ self._ensure_loop(task)
return task
+ def _ensure_loop(self, task: ScheduledTask) -> None:
+ """为任务启动运行循环(幂等:已在运行则跳过)。"""
+ if task._task is None or task._task.done():
+ task._task = asyncio.create_task(task._run_loop())
+
def get(self, task_id: str) -> ScheduledTask | None:
return self._tasks.get(task_id)
@@ -174,14 +230,18 @@ class Scheduler:
return list(self._tasks.values())
def remove(self, task_id: str) -> None:
- self._tasks.pop(task_id, None)
+ task = self._tasks.pop(task_id, None)
+ if task is not None and task._task is not None:
+ task._task.cancel()
async def start(self) -> None:
+ self._running = True
for task in self._tasks.values():
- task._task = asyncio.create_task(task._run_loop())
+ self._ensure_loop(task)
logger.info("定时调度器已启动, 共 %d 个任务", len(self._tasks))
async def stop(self) -> None:
+ self._running = False
for task in self._tasks.values():
if task._task:
task._task.cancel()
diff --git a/tests/test_scheduler_registration.py b/tests/test_scheduler_registration.py
new file mode 100644
index 0000000..e7604c9
--- /dev/null
+++ b/tests/test_scheduler_registration.py
@@ -0,0 +1,190 @@
+"""C1 回归测试: 定时任务调度器「静默失效」缺陷。
+
+原始缺陷(全量审查 C1):
+ 1. `Scheduler.register()` 的第 2 参是 cron 表达式,但 app.py:52 /
+ schedules.py:95,119 三处调用传的都是任务类型字符串("events" 等)。
+ `croniter("events")` 抛异常被 `_calc_next` 的 except 吞掉 →
+ next_run 永远为 None → 所有定时任务从不触发且无任何报错。
+ 2. 运行期通过 API 新建的任务只进了 _tasks 字典,从未
+ asyncio.create_task 启动 _run_loop(只有 Scheduler.start() 会建循环)。
+
+这些用例不依赖数据库,直接对调度器类做行为断言。
+"""
+from __future__ import annotations
+
+import asyncio
+
+import pytest
+
+from src.core.scheduler import ScheduledTask, Scheduler
+
+
+class TestCronValidation:
+ """非法 cron 必须显式报错,不能静默变成永不触发。"""
+
+ def test_nonevent_cron_raises(self):
+ """把任务类型字符串当 cron 传(原缺陷) → 构造时立即抛 ValueError。"""
+ async def _noop() -> None: ...
+ with pytest.raises(ValueError, match="cron"):
+ ScheduledTask("daily-events", "events", _noop)
+
+ def test_valid_cron_accepted(self):
+ """合法 cron 正常构造,且 next_run 被计算出来(非 None)。"""
+ async def _noop() -> None: ...
+ task = ScheduledTask("daily-events", "0 8 * * *", _noop)
+ assert task.next_run is not None
+
+ def test_next_run_is_in_future(self):
+ """next_run 必须落在未来,否则 _run_loop 会立刻误触发。"""
+ from datetime import datetime
+ async def _noop() -> None: ...
+ task = ScheduledTask("t", "*/30 * * * *", _noop)
+ assert task.next_run > datetime.now()
+
+ def test_croniter_receives_called_datetime(self):
+ """回归: `croniter(expr, datetime.now)`(未调用)会抛
+ TypeError: 'builtin_function_or_method' object cannot be interpreted
+ as an integer —— 这是比「传错参数」更深一层的静默失效根源。
+ 断言 next_run 是真实 datetime,而非 None。"""
+ from datetime import datetime
+ async def _noop() -> None: ...
+ task = ScheduledTask("t", "0 8 * * *", _noop)
+ assert isinstance(task.next_run, datetime), (
+ f"next_run 应为 datetime,实际 {task.next_run!r} —— "
+ "croniter 的 start_time 未正确传入"
+ )
+
+
+class TestRunLoopActuallyFires:
+ """注册后任务必须真的在到期时被执行(端到端行为,不 mock 内部)。"""
+
+ async def test_register_then_start_runs_task(self):
+ """cron 到点后任务函数被调用一次。用 1 秒粒度的 * * * * * 加速验证。"""
+ calls: list[str] = []
+
+ async def _job() -> None:
+ calls.append("ran")
+
+ sched = Scheduler()
+ sched.register("t1", "* * * * *", _job, enabled=True)
+ await sched.start()
+ try:
+ # _run_loop 用 min(wait, 60) 分段睡;下一分钟边界最多 60 秒。
+ # 为让测试可跑,直接把 next_run 拨到过去,触发一次执行。
+ task = sched.get("t1")
+ from datetime import datetime, timedelta
+ task.next_run = datetime.now() - timedelta(seconds=1)
+ for _ in range(40):
+ if calls:
+ break
+ await asyncio.sleep(0.05)
+ finally:
+ await sched.stop()
+ assert calls == ["ran"], "注册并 start 后,到期任务未被执行"
+
+ async def test_register_after_start_also_runs(self):
+ """运行期(已 start 之后)新 register 的任务也必须被启动(原缺陷 2)。"""
+ calls: list[str] = []
+
+ async def _job() -> None:
+ calls.append("late")
+
+ sched = Scheduler()
+ await sched.start()
+ try:
+ sched.register("late-task", "* * * * *", _job, enabled=True)
+ from datetime import datetime, timedelta
+ sched.get("late-task").next_run = datetime.now() - timedelta(seconds=1)
+ for _ in range(40):
+ if calls:
+ break
+ await asyncio.sleep(0.05)
+ finally:
+ await sched.stop()
+ assert calls == ["late"], "start 之后注册的任务未被启动循环"
+
+ async def test_disabled_task_does_not_run(self):
+ """enabled=False 的任务不执行。"""
+ calls: list[str] = []
+
+ async def _job() -> None:
+ calls.append("nope")
+
+ sched = Scheduler()
+ sched.register("off", "* * * * *", _job, enabled=False)
+ await sched.start()
+ try:
+ from datetime import datetime, timedelta
+ sched.get("off").next_run = datetime.now() - timedelta(seconds=1)
+ await asyncio.sleep(0.3)
+ finally:
+ await sched.stop()
+ assert calls == []
+
+
+class TestTaskFailureIsolation:
+ """单个任务抛异常不能杀掉循环(否则一次失败永久停摆)。"""
+
+ async def test_exception_does_not_kill_loop(self):
+ """任务抛异常后,循环仍存活并能在下一次到期时继续执行。"""
+ runs: list[int] = []
+
+ async def _boom() -> None:
+ runs.append(1)
+ if len(runs) == 1:
+ raise RuntimeError("boom")
+
+ sched = Scheduler()
+ sched.register("flaky", "* * * * *", _boom, enabled=True)
+ await sched.start()
+ try:
+ from datetime import datetime, timedelta
+ task = sched.get("flaky")
+ for _ in range(40):
+ if len(runs) >= 2:
+ break
+ task.next_run = datetime.now() - timedelta(seconds=1)
+ await asyncio.sleep(0.05)
+ assert len(runs) >= 2, "任务首次抛异常后循环未继续"
+ assert task._task is not None and not task._task.done(), (
+ "循环在任务抛异常后已终止"
+ )
+ finally:
+ await sched.stop()
+
+ async def test_next_run_change_takes_effect_promptly(self):
+ """运行期把 next_run 改到过去,循环须在 SLEEP_TICK 内感知并触发。
+
+ 回归: 原实现 sleep(min(wait_seconds, 60)),当 wait_seconds<60 时
+ 会一次性睡满 wait_seconds,导致运行期通过 API 更新 cron 后
+ 最长 60s 不生效。
+ """
+ calls: list[str] = []
+
+ async def _job() -> None:
+ calls.append("ran")
+
+ import src.core.scheduler as sched_mod
+ old_tick = sched_mod.ScheduledTask.SLEEP_TICK
+ sched_mod.ScheduledTask.SLEEP_TICK = 0.05
+ try:
+ sched = Scheduler()
+ sched.register("tick", "* * * * *", _job, enabled=True)
+ await sched.start()
+ try:
+ from datetime import datetime, timedelta
+ # 先等循环进入 sleep(wait 接近 60s 的最坏情况)
+ await asyncio.sleep(0.15)
+ sched.get("tick").next_run = datetime.now() - timedelta(seconds=1)
+ # 短 tick 应在 ~0.15s 内感知;给 1s 容差
+ for _ in range(40):
+ if calls:
+ break
+ await asyncio.sleep(0.05)
+ finally:
+ await sched.stop()
+ finally:
+ sched_mod.ScheduledTask.SLEEP_TICK = old_tick
+ assert calls == ["ran"], (
+ "next_run 变更后循环未在 tick 内响应(疑似一次性睡满 wait_seconds)"
+ )