Spaces:
Sleeping
Sleeping
File size: 6,239 Bytes
bdef324 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 | # scheduler.py — Multiple-schedule forecast manager
# Ported from Daily Scrapper architecture with forecast-specific job
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.triggers.cron import CronTrigger
from logger import get_app_logger
from datetime import datetime
import uuid
try:
from pytz import timezone as pytz_timezone
wib_tz = pytz_timezone("Asia/Jakarta")
except ImportError:
import zoneinfo
wib_tz = zoneinfo.ZoneInfo("Asia/Jakarta")
logger = get_app_logger("scheduler")
class ForecastScheduler:
"""Manages multiple forecast schedules with APScheduler."""
def __init__(self):
self.scheduler = BackgroundScheduler(timezone=wib_tz)
self.is_running = False
self.last_run = None
self.last_result = None
self.schedules = {} # id -> {"cron": str, "label": str}
self._forecast_job_fn = None # Set by main.py
def set_forecast_job(self, fn):
"""Register the forecast function to be called by scheduled jobs."""
self._forecast_job_fn = fn
def _run_job(self):
"""Execute the registered forecast job."""
logger.info("⏰ Menjalankan forecast terjadwal otomatis...")
try:
if self._forecast_job_fn:
self._forecast_job_fn()
self.last_run = datetime.now()
self.last_result = {
"success": True,
"timestamp": self.last_run.isoformat(),
}
logger.info("✅ Forecast terjadwal selesai.")
except Exception as e:
self.last_result = {
"success": False,
"error": str(e),
"timestamp": datetime.now().isoformat(),
}
logger.error(f"❌ Forecast terjadwal gagal: {e}")
def start_all(self, default_schedules: str = ""):
"""Start the scheduler engine. Optionally load default schedules."""
if not self.is_running:
self.scheduler = BackgroundScheduler(timezone=wib_tz)
if self._forecast_job_fn:
# Re-register all existing schedules
for jid, info in list(self.schedules.items()):
try:
trigger = CronTrigger.from_crontab(info["cron"], timezone=wib_tz)
self.scheduler.add_job(
self._run_job, trigger, id=jid, replace_existing=True
)
except Exception:
pass
self.scheduler.start()
self.is_running = True
logger.info("▶️ Mesin Scheduler Forecast telah diaktifkan.")
# Load defaults if no schedules exist
if not self.schedules and default_schedules:
for cron in [s.strip() for s in default_schedules.split(",") if s.strip()]:
self.add_schedule(cron, label="Jadwal Default")
def stop_all(self):
"""Stop the scheduler engine and clear all jobs."""
if self.is_running:
try:
self.scheduler.shutdown(wait=False)
except Exception:
pass
self.scheduler = BackgroundScheduler(timezone=wib_tz)
self.is_running = False
self.schedules.clear()
logger.info("⏹️ Mesin Scheduler Forecast telah dimatikan.")
def add_schedule(self, cron_expression: str, label: str = "") -> str:
"""Add a new cron schedule. Returns job_id."""
job_id = str(uuid.uuid4())[:8]
try:
trigger = CronTrigger.from_crontab(cron_expression, timezone=wib_tz)
self.scheduler.add_job(
self._run_job, trigger, id=job_id, replace_existing=True
)
if not label:
label = f"Jadwal {cron_expression}"
self.schedules[job_id] = {"cron": cron_expression, "label": label}
logger.info(f"📅 Jadwal baru ditambahkan: {cron_expression} ({label})")
# Auto-start if not running
if not self.is_running:
self.scheduler.start()
self.is_running = True
logger.info("▶️ Scheduler otomatis menyala karena ada jadwal baru.")
return job_id
except Exception as e:
logger.error(f"❌ Gagal menambahkan jadwal cron '{cron_expression}': {e}")
raise ValueError(f"Format Cron tidak valid: {e}")
def remove_schedule(self, job_id: str):
"""Remove a schedule by job_id."""
if job_id in self.schedules:
try:
self.scheduler.remove_job(job_id)
except Exception:
pass
deleted = self.schedules.pop(job_id)
logger.info(f"🗑️ Jadwal dihapus: {deleted['cron']} ({deleted['label']})")
else:
raise KeyError("ID Jadwal tidak ditemukan")
def list_schedules(self):
"""List all active schedules with next_run info."""
result = []
jobs_map = {job.id: job for job in self.scheduler.get_jobs()}
for jid, info in self.schedules.items():
job_obj = jobs_map.get(jid)
next_run = (
job_obj.next_run_time.isoformat()
if (job_obj and job_obj.next_run_time)
else None
)
result.append({
"id": jid,
"cron_expression": info["cron"],
"label": info["label"],
"next_run": next_run,
})
return result
def get_status(self):
"""Get overall scheduler status."""
jobs = self.scheduler.get_jobs() if self.is_running else []
next_runs = [j.next_run_time for j in jobs if j.next_run_time]
next_run_overall = min(next_runs).isoformat() if next_runs else None
return {
"is_running": self.is_running,
"last_run": self.last_run.isoformat() if self.last_run else None,
"last_result": self.last_result,
"next_run": next_run_overall,
"active_schedules_count": len(self.schedules),
}
# Singleton instance
forecast_scheduler = ForecastScheduler()
|