From 8c87b9d0249cf5c5e62fb223358552a84b9200a9 Mon Sep 17 00:00:00 2001 From: shangfangjian Date: Mon, 21 Sep 2026 17:04:04 +0800 Subject: [PATCH] =?UTF-8?q?fix(critical):=20=E4=BF=AE=E5=A4=8D=E8=B0=83?= =?UTF-8?q?=E5=BA=A6=E5=99=A8=E9=9D=99=E9=BB=98=E5=A4=B1=E6=95=88=E3=80=81?= =?UTF-8?q?=E6=95=B0=E6=8D=AE=E6=BA=90=E7=BC=BA=E5=A4=B1=E4=B8=8E=E5=89=8D?= =?UTF-8?q?=E7=AB=AF=E7=B1=BB=E5=9E=8B=E9=94=99=E8=AF=AF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit C1 定时任务静默失效(4 层缺陷): - scheduler.py: croniter(expr, datetime.now) 未调用 now(), 导致 TypeError 被吞、next_run 恒为 None(最深一层) - scheduler.py: _calc_next 非法 cron 静默吞异常 -> 改为构造期抛 ValueError - scheduler.py: _run_loop 一次异常即永久停摆 -> 增加异常隔离 - scheduler.py: sleep(min(wait,60)) 使运行期 cron 变更最长 60s 才生效 -> 改为固定 SLEEP_TICK 轮询 - app.py / schedules.py: 3 处 register() 第 2 参误传 task 类型名 C3 EvalPage.tsx: 引用未导入的 EmptyState -> 改为已导入的 EmptyText (tsc --noEmit 由 1 error 变为 0 error) 新增 tests/test_scheduler_registration.py: 9 项行为回归测试 --- frontend/src/admin/pages/EvalPage.tsx | 2 +- src/api/app.py | 2 +- src/api/routes/schedules.py | 8 +- src/core/scheduler.py | 118 ++++++++++++---- tests/test_scheduler_registration.py | 190 ++++++++++++++++++++++++++ 5 files changed, 285 insertions(+), 35 deletions(-) create mode 100644 tests/test_scheduler_registration.py 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)" + )