feat: 调度简单设置(频率选择器)+ 运行窗口 + 下次运行列
This commit is contained in:
+72
-1
@@ -21,6 +21,7 @@ import json
|
||||
import time
|
||||
import uuid
|
||||
import threading
|
||||
from datetime import datetime
|
||||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||
|
||||
from apscheduler.schedulers.background import BackgroundScheduler
|
||||
@@ -45,6 +46,61 @@ _log = get_logger("core.tm")
|
||||
_START_STAGGER_SEC = 0.2
|
||||
|
||||
|
||||
def _in_run_window(schedule, now=None):
|
||||
"""是否在当前运行窗口内。
|
||||
|
||||
schedule.window = {"start": "HH:MM", "end": "HH:MM"},每天重复;
|
||||
窗口外(定时触发 + 手动执行)任务不会启动。未配置/非法配置视为不限制。
|
||||
支持跨午夜(如 21:00-09:00 = 晚上 9 点运行到次日早 9 点)。
|
||||
"""
|
||||
win = (schedule or {}).get("window") or {}
|
||||
start, end = win.get("start", ""), win.get("end", "")
|
||||
if not start or not end:
|
||||
return True
|
||||
try:
|
||||
cur = (now or datetime.now()).hour * 60 + (now or datetime.now()).minute
|
||||
s = int(start.split(":")[0]) * 60 + int(start.split(":")[1])
|
||||
e = int(end.split(":")[0]) * 60 + int(end.split(":")[1])
|
||||
except (ValueError, AttributeError, IndexError):
|
||||
return True # 配置非法按不限制处理
|
||||
if s == e:
|
||||
return True # 起止相同视为不限制
|
||||
if s < e:
|
||||
return s <= cur < e
|
||||
return cur >= s or cur < e # 跨午夜
|
||||
|
||||
|
||||
def _next_run_time(schedule, now=None):
|
||||
"""任务下次真正执行的时间(考虑运行窗口)。
|
||||
|
||||
用 CronTrigger.get_next_fire_time 从当前时间向后找触发点,
|
||||
跳过运行窗口外的触发点(最多找 200 次防死循环)。
|
||||
返回 datetime 或 None(无 cron / 配置非法)。
|
||||
"""
|
||||
sched = schedule or {}
|
||||
if sched.get("mode") not in ("cron", "cron_stop"):
|
||||
return None
|
||||
cron = sched.get("cron", "")
|
||||
if not cron:
|
||||
return None
|
||||
try:
|
||||
trigger = CronTrigger.from_crontab(cron)
|
||||
except Exception:
|
||||
return None
|
||||
now = now or datetime.now()
|
||||
win = sched.get("window") or {}
|
||||
if not win.get("start") or not win.get("end"):
|
||||
return trigger.get_next_fire_time(None, now)
|
||||
fire = trigger.get_next_fire_time(None, now)
|
||||
for _ in range(200):
|
||||
if fire is None:
|
||||
return None
|
||||
if _in_run_window({"window": win}, fire):
|
||||
return fire
|
||||
fire = trigger.get_next_fire_time(fire, fire)
|
||||
return None
|
||||
|
||||
|
||||
# ================== 设备分组 ==================
|
||||
class DeviceGroup:
|
||||
def __init__(self, name, serials=None, description=""):
|
||||
@@ -504,6 +560,11 @@ class TaskManager:
|
||||
job = self.jobs.get(job_id)
|
||||
if not job:
|
||||
return
|
||||
if not _in_run_window(job.schedule):
|
||||
win = job.schedule.get("window") or {}
|
||||
_log.info(f"定时触发 {job.name}({job.id}) 跳过:当前不在运行窗口内 "
|
||||
f"({win.get('start','')}-{win.get('end','')})")
|
||||
return
|
||||
_log.info(f"定时触发: {job.name}({job.id})")
|
||||
self._run_job(job)
|
||||
|
||||
@@ -526,11 +587,21 @@ class TaskManager:
|
||||
stopped.append(serial)
|
||||
_log.info(f"定时停止 {job.name}({job.id}): 停止 {len(stopped)} 台设备 {stopped}")
|
||||
|
||||
def next_run_of(self, job):
|
||||
"""任务下次真正执行的时间(考虑运行窗口)。停用/无 cron 返回 None。"""
|
||||
if not job or not job.enabled:
|
||||
return None
|
||||
return _next_run_time(job.schedule)
|
||||
|
||||
def run_job_now(self, job_id):
|
||||
"""立即执行任务(手动触发)。"""
|
||||
"""立即执行任务(手动触发)。运行窗口外拒绝启动。"""
|
||||
job = self.jobs.get(job_id)
|
||||
if not job:
|
||||
return {"ok": False, "error": "任务不存在"}
|
||||
if not _in_run_window(job.schedule):
|
||||
win = job.schedule.get("window") or {}
|
||||
return {"ok": False,
|
||||
"error": f"当前不在运行窗口内({win.get('start', '')}-{win.get('end', '')}),任务未启动"}
|
||||
# 在独立线程跑,不阻塞调用方
|
||||
t = threading.Thread(target=self._run_job, args=(job,), daemon=True)
|
||||
t.start()
|
||||
|
||||
Reference in New Issue
Block a user