scheduler.py 4.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117
  1. from apscheduler.schedulers.background import BackgroundScheduler
  2. from apscheduler.triggers.cron import CronTrigger
  3. _flask_app = None
  4. scheduler = BackgroundScheduler(
  5. # misfire_grace_time généreux : par défaut (1s), un job planifié à la même
  6. # seconde que d'autres (ex: plusieurs "0 3 * * *") peut être traité trop
  7. # tard par l'ordonnanceur et est alors sauté silencieusement, sans trace.
  8. job_defaults={"coalesce": True, "max_instances": 1, "misfire_grace_time": 3600},
  9. timezone="UTC",
  10. )
  11. def init_scheduler(flask_app):
  12. global _flask_app
  13. _flask_app = flask_app
  14. # Nettoyage immédiat au démarrage : tout run "running" en DB est forcément
  15. # un vestige d'un process précédent (app redémarrée pendant un backup).
  16. with flask_app.app_context():
  17. from db import db, Run
  18. import datetime as _dt
  19. stale = Run.query.filter_by(status="running").all()
  20. now = _dt.datetime.utcnow()
  21. for run in stale:
  22. run.status = "error"
  23. run.log_text = (run.log_text or "") + "\n[interrompu] Run interrompu par un redémarrage de l'application."
  24. run.finished_at = now
  25. if stale:
  26. db.session.commit()
  27. if not scheduler.running:
  28. scheduler.start()
  29. # Filet de sécurité : marque en erreur tout run resté bloqué > 6h
  30. scheduler.add_job(
  31. func=_cleanup_stuck_runs,
  32. trigger="interval",
  33. hours=1,
  34. id="cleanup_stuck_runs",
  35. replace_existing=True,
  36. )
  37. def _cleanup_stuck_runs():
  38. import datetime as _dt
  39. with _flask_app.app_context():
  40. from db import db, Run
  41. cutoff = _dt.datetime.utcnow() - _dt.timedelta(hours=6)
  42. stuck = Run.query.filter(
  43. Run.status == "running",
  44. Run.started_at < cutoff,
  45. ).all()
  46. now = _dt.datetime.utcnow()
  47. for run in stuck:
  48. duration = now - run.started_at
  49. hours = int(duration.total_seconds() // 3600)
  50. minutes = int((duration.total_seconds() % 3600) // 60)
  51. run.status = "error"
  52. run.finished_at = now
  53. if run.archive_name:
  54. # L'archive a été créée localement ; le timeout est survenu pendant le transfert
  55. detail = (
  56. f" Archive locale présente ({run.archive_name})."
  57. " Le processus de transfert était probablement encore en cours."
  58. )
  59. else:
  60. # Le backup lui-même n'a pas terminé dans les 6h
  61. detail = " Aucune archive locale détectée : la création du backup n'a pas abouti dans les délais."
  62. run.log_text = (run.log_text or "") + (
  63. f"\n\n[timeout] Run bloqué depuis {hours}h{minutes:02d} — "
  64. f"marqué en erreur par le nettoyage automatique.{detail}"
  65. )
  66. if stuck:
  67. db.session.commit()
  68. def _execute_job(job_id):
  69. with _flask_app.app_context():
  70. from jobs.ynh_backup import execute_job
  71. execute_job(job_id)
  72. def schedule_job(job):
  73. import logging
  74. job_key = f"job_{job.id}"
  75. if not job.cron_expr:
  76. return # job manuel uniquement, pas de planification APScheduler
  77. try:
  78. trigger = CronTrigger.from_crontab(job.cron_expr)
  79. except Exception:
  80. logging.warning(f"Job #{job.id} « {job.name} » : expression cron invalide « {job.cron_expr} » — job non planifié.")
  81. return
  82. if scheduler.get_job(job_key):
  83. scheduler.reschedule_job(job_key, trigger=trigger)
  84. else:
  85. scheduler.add_job(
  86. func=_execute_job,
  87. trigger=trigger,
  88. id=job_key,
  89. kwargs={"job_id": job.id},
  90. replace_existing=True,
  91. )
  92. def remove_job(job_id):
  93. job_key = f"job_{job_id}"
  94. if scheduler.get_job(job_key):
  95. scheduler.remove_job(job_key)
  96. def get_next_run(job_id):
  97. job_key = f"job_{job_id}"
  98. apsjob = scheduler.get_job(job_key)
  99. if apsjob and apsjob.next_run_time:
  100. return apsjob.next_run_time
  101. return None