scheduled_trigger.py 9.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281
  1. # coding=utf-8
  2. from django.db.models import QuerySet
  3. from common.utils.logger import maxkb_logger
  4. from ops import celery_app
  5. from trigger.handler.base_trigger import BaseTrigger
  6. from trigger.models import TriggerTask
  7. def _parse_hhmm(value: str) -> tuple[int, int]:
  8. hour_str, minute_str = (value or "").split(":")
  9. hour = int(hour_str)
  10. minute = int(minute_str)
  11. if not (0 <= hour <= 23 and 0 <= minute <= 59):
  12. raise ValueError("hour/minute out of range")
  13. return hour, minute
  14. def _weekday_to_cron(d: int | str) -> str:
  15. mapping = {1: "mon", 2: "tue", 3: "wed", 4: "thu", 5: "fri", 6: "sat", 7: "sun", 0: "sun"}
  16. di = int(d)
  17. if di not in mapping:
  18. raise ValueError("invalid weekday")
  19. return mapping[di]
  20. def _get_active_trigger_tasks(trigger_id: str) -> list[dict]:
  21. return list(
  22. QuerySet(TriggerTask)
  23. .filter(trigger_id=trigger_id, is_active=True)
  24. .values("id", "source_type", "source_id", "parameter", "trigger")
  25. )
  26. def _deploy_daily(trigger: dict, trigger_tasks: list[dict], setting: dict, trigger_id: str) -> None:
  27. from common.job import scheduler
  28. times = setting.get("time") or []
  29. for t in times:
  30. try:
  31. hour, minute = _parse_hhmm(t)
  32. except Exception:
  33. maxkb_logger.warning(f"invalid time={t}, trigger_id={trigger_id}")
  34. continue
  35. for task in trigger_tasks:
  36. job_id = f"trigger:{trigger_id}:task:{task['id']}:daily:{hour:02d}{minute:02d}"
  37. scheduler.add_job(
  38. ScheduledTrigger.execute,
  39. trigger="cron",
  40. hour=str(hour),
  41. minute=str(minute),
  42. id=job_id,
  43. kwargs={"trigger": trigger, "trigger_task": task},
  44. replace_existing=True,
  45. misfire_grace_time=60,
  46. max_instances=1,
  47. )
  48. def _deploy_weekly(trigger: dict, trigger_tasks: list[dict], setting: dict, trigger_id: str) -> None:
  49. from common.job import scheduler
  50. times = setting.get("time") or []
  51. days = setting.get("days") or []
  52. if not times or not days:
  53. maxkb_logger.warning(f"empty weekly setting, trigger_id={trigger_id}")
  54. return
  55. for d in days:
  56. try:
  57. dow = _weekday_to_cron(d)
  58. except Exception:
  59. maxkb_logger.warning(f"invalid weekday={d}, trigger_id={trigger_id}")
  60. continue
  61. for t in times:
  62. try:
  63. hour, minute = _parse_hhmm(t)
  64. except Exception:
  65. maxkb_logger.warning(f"invalid time={t}, trigger_id={trigger_id}")
  66. continue
  67. for task in trigger_tasks:
  68. job_id = f"trigger:{trigger_id}:task:{task['id']}:weekly:{dow}:{hour:02d}{minute:02d}"
  69. scheduler.add_job(
  70. ScheduledTrigger.execute,
  71. trigger="cron",
  72. day_of_week=dow,
  73. hour=str(hour),
  74. minute=str(minute),
  75. id=job_id,
  76. kwargs={"trigger": trigger, "trigger_task": task},
  77. replace_existing=True,
  78. misfire_grace_time=60,
  79. max_instances=1,
  80. )
  81. def _deploy_monthly(trigger: dict, trigger_tasks: list[dict], setting: dict, trigger_id: str) -> None:
  82. from common.job import scheduler
  83. times = setting.get("time") or []
  84. days = setting.get("days") or []
  85. if not times or not days:
  86. maxkb_logger.warning(f"empty monthly setting, trigger_id={trigger_id}")
  87. return
  88. for d in days:
  89. try:
  90. dom = int(d)
  91. if not (1 <= dom <= 31):
  92. raise ValueError("invalid day of month")
  93. except Exception:
  94. maxkb_logger.warning(f"invalid day={d}, trigger_id={trigger_id}")
  95. continue
  96. for t in times:
  97. try:
  98. hour, minute = _parse_hhmm(t)
  99. except Exception:
  100. maxkb_logger.warning(f"invalid time={t}, trigger_id={trigger_id}")
  101. continue
  102. for task in trigger_tasks:
  103. job_id = f"trigger:{trigger_id}:task:{task['id']}:monthly:{dom:02d}:{hour:02d}{minute:02d}"
  104. scheduler.add_job(
  105. ScheduledTrigger.execute,
  106. trigger="cron",
  107. day=str(dom),
  108. hour=str(hour),
  109. minute=str(minute),
  110. id=job_id,
  111. kwargs={"trigger": trigger, "trigger_task": task},
  112. replace_existing=True,
  113. misfire_grace_time=60,
  114. max_instances=1,
  115. )
  116. def _deploy_cron(trigger: dict, trigger_tasks: list[dict], setting: dict, trigger_id: str) -> None:
  117. from common.job import scheduler
  118. from apscheduler.triggers.cron import CronTrigger
  119. cron_expression = setting.get('cron_expression')
  120. if not cron_expression:
  121. maxkb_logger.warning(f"empty cron_expression, trigger_id={trigger_id}")
  122. return
  123. try:
  124. cron_trigger = CronTrigger.from_crontab(cron_expression.strip())
  125. except ValueError:
  126. maxkb_logger.warning(f"invalid cron_expression={cron_expression}, trigger_id={trigger_id}")
  127. return
  128. for task in trigger_tasks:
  129. job_id = f"trigger:{trigger_id}:task:{task['id']}:cron:{cron_expression.strip()}"
  130. scheduler.add_job(
  131. ScheduledTrigger.execute,
  132. trigger=cron_trigger,
  133. id=job_id,
  134. kwargs={"trigger": trigger, "trigger_task": task},
  135. replace_existing=True,
  136. misfire_grace_time=60,
  137. max_instances=1,
  138. )
  139. def _deploy_interval(trigger: dict, trigger_tasks: list[dict], setting: dict, trigger_id: str) -> None:
  140. from common.job import scheduler
  141. unit = (setting.get("interval_unit") or "").strip()
  142. value = setting.get("interval_value")
  143. try:
  144. value_i = int(value)
  145. if value_i <= 0:
  146. raise ValueError("interval_value must be positive")
  147. except Exception:
  148. maxkb_logger.warning(f"invalid interval_value={value}, trigger_id={trigger_id}")
  149. return
  150. if unit not in {"seconds", "minutes", "hours", "days"}:
  151. maxkb_logger.warning(f"invalid interval_unit={unit}, trigger_id={trigger_id}")
  152. return
  153. for task in trigger_tasks:
  154. job_id = f"trigger:{trigger_id}:task:{task['id']}:interval:{unit}:{value_i}"
  155. scheduler.add_job(
  156. ScheduledTrigger.execute,
  157. trigger="interval",
  158. id=job_id,
  159. kwargs={"trigger": trigger, "trigger_task": task},
  160. replace_existing=True,
  161. misfire_grace_time=60,
  162. max_instances=1,
  163. **{unit: value_i},
  164. )
  165. @celery_app.task(name="celery:undeploy_scheduled_trigger")
  166. def _remove_trigger_jobs(trigger_id: str) -> None:
  167. from common.job import scheduler
  168. prefix = f"trigger:{trigger_id}:"
  169. for job in scheduler.get_jobs():
  170. if getattr(job, "id", "").startswith(prefix):
  171. try:
  172. job.remove()
  173. except Exception as e:
  174. maxkb_logger.warning(f"remove job failed, job_id={job.id}, err={e}")
  175. @celery_app.task(name="celery:deploy_scheduled_trigger")
  176. def deploy_scheduled_trigger(trigger: dict, trigger_tasks: list[dict], setting: dict, schedule_type: str) -> None:
  177. _remove_trigger_jobs(trigger["id"])
  178. deployers = {
  179. "daily": _deploy_daily,
  180. "weekly": _deploy_weekly,
  181. "monthly": _deploy_monthly,
  182. "interval": _deploy_interval,
  183. 'cron': _deploy_cron
  184. }
  185. fn = deployers.get(schedule_type)
  186. if not fn:
  187. maxkb_logger.warning(f"unsupported schedule_type={schedule_type}, trigger_id={trigger['id']}")
  188. return
  189. fn(trigger, trigger_tasks, setting, trigger["id"])
  190. class ScheduledTrigger(BaseTrigger):
  191. """
  192. 定时任务触发器
  193. """
  194. @staticmethod
  195. def execute(trigger, **kwargs):
  196. trigger_task = kwargs.get("trigger_task")
  197. if not trigger_task:
  198. maxkb_logger.warning(f"unsupported task={trigger_task}")
  199. return
  200. source_type = trigger_task["source_type"]
  201. if source_type == "APPLICATION":
  202. from trigger.handler.impl.task.application_task import ApplicationTask
  203. ApplicationTask.execute(trigger_task, **kwargs)
  204. elif source_type == "TOOL":
  205. from trigger.handler.impl.task.tool_task import ToolTask
  206. ToolTask.execute(trigger_task, **kwargs)
  207. else:
  208. maxkb_logger.warning(f"unsupported source_type={source_type}, task_id={trigger_task['id']}")
  209. def support(self, trigger, **kwargs):
  210. return trigger.get("trigger_type") == "SCHEDULED"
  211. def deploy(self, trigger, **kwargs):
  212. trigger_id = str(trigger["id"])
  213. setting = trigger.get("trigger_setting") or {}
  214. schedule_type = setting.get("schedule_type")
  215. if not trigger.get("is_active", True):
  216. self.undeploy(trigger, **kwargs)
  217. return
  218. if trigger.get("trigger_type") != "SCHEDULED":
  219. self.undeploy(trigger, **kwargs)
  220. return
  221. trigger_tasks = _get_active_trigger_tasks(trigger["id"])
  222. if not trigger_tasks:
  223. maxkb_logger.warning(f"no active trigger_tasks, trigger_id={trigger_id}")
  224. self.undeploy(trigger, **kwargs)
  225. return
  226. deploy_scheduled_trigger.delay(trigger, trigger_tasks, setting, schedule_type)
  227. def undeploy(self, trigger, **kwargs):
  228. trigger_id = str(trigger["id"])
  229. _remove_trigger_jobs.delay(trigger_id)