""" Scheduler Daemon for Python Scheduler Main daemon process that manages job scheduling and execution using APScheduler """ import time import signal import sys from typing import Optional from datetime import datetime from apscheduler.schedulers.background import BackgroundScheduler from apscheduler.triggers.cron import CronTrigger from apscheduler.triggers.interval import IntervalTrigger from apscheduler.triggers.date import DateTrigger from loguru import logger from core.database import SchedulerDatabase from core.config import SchedulerConfig from core.logger import setup_logging from core.cron_utils import convert_unix_to_iso_dayofweek from daemon.job_executor import JobExecutor from daemon.notification_manager import NotificationManager from daemon.dependency_manager import DependencyManager from daemon.relay_handler import process_pending_requests class SchedulerDaemon: """Main scheduler daemon using APScheduler""" def __init__(self): # Setup logging first setup_logging('daemon') logger.info("Initializing Scheduler Daemon") # Load configuration self.config = SchedulerConfig() # Initialize database (con retry se B non รจ raggiungibile) self.db = self._init_database_with_retry() # Initialize notification manager self.notification_manager = NotificationManager(self.db, self.config) # Initialize job executor with notification manager self.executor = JobExecutor(self.db, self.config, self.notification_manager) # Initialize dependency manager self.dependency_manager = DependencyManager(self.db, self.executor) # Initialize APScheduler self.scheduler = BackgroundScheduler( timezone='Europe/Rome', # or use your timezone job_defaults={ 'coalesce': True, # Combine multiple missed runs into one 'max_instances': 1, # Prevent overlapping executions 'misfire_grace_time': 300 # 5 minutes grace period } ) # Shutdown flag self.running = False # Setup signal handlers signal.signal(signal.SIGINT, self._signal_handler) signal.signal(signal.SIGTERM, self._signal_handler) logger.info("Scheduler Daemon initialized") def _init_database_with_retry(self) -> SchedulerDatabase: """Inizializza il DB con retry โ€” gestisce B non raggiungibile all'avvio (es. VPN non ancora attiva)""" db_path = self.config.database_path retry_interval = self.config.get('database.retry_interval_seconds', 30) max_attempts = self.config.get('database.retry_max_attempts', 60) attempt = 0 while True: try: db = SchedulerDatabase(db_path) if attempt > 0: logger.info(f"Database connesso dopo {attempt} tentativo/i.") return db except Exception as e: attempt += 1 if max_attempts and attempt >= max_attempts: logger.error(f"Database non raggiungibile dopo {attempt} tentativi. Arresto.") raise logger.warning( f"Database non raggiungibile ({db_path}): {e}. " f"Tentativo {attempt}/{max_attempts if max_attempts else 'โˆž'} โ€” " f"nuovo tentativo tra {retry_interval}s..." ) time.sleep(retry_interval) def _signal_handler(self, signum, frame): """Handle shutdown signals""" logger.info(f"Received signal {signum}, shutting down...") self.shutdown() def start(self): """Start the scheduler daemon""" logger.info("Starting Scheduler Daemon") try: # Load all enabled schedules from database self._load_schedules() # Start APScheduler self.scheduler.start() self.running = True logger.info("Scheduler Daemon started successfully") logger.info(f"Loaded {len(self.scheduler.get_jobs())} schedules") # Main loop - keep daemon alive try: while self.running: time.sleep(self.config.get('daemon.check_interval_seconds', 30)) # Optional: Periodic health check or cleanup self._health_check() except KeyboardInterrupt: logger.info("Keyboard interrupt received") except Exception as e: logger.exception(f"Error starting daemon: {e}") raise finally: self.shutdown() def shutdown(self): """Shutdown the daemon gracefully""" if not self.running: return logger.info("Shutting down Scheduler Daemon") self.running = False # Stop scheduler if self.scheduler.running: self.scheduler.shutdown(wait=True) logger.info("APScheduler stopped") # Wait for active jobs to complete (with timeout) active = self.executor.get_active_executions() if active: logger.info(f"Waiting for {len(active)} active jobs to complete...") timeout = 60 # Wait max 60 seconds start = time.time() while active and (time.time() - start) < timeout: time.sleep(1) active = self.executor.get_active_executions() if active: logger.warning(f"{len(active)} jobs still running after timeout") logger.info("Scheduler Daemon stopped") def _load_schedules(self): """Load all enabled schedules from database and schedule them""" schedules = self.db.get_all_schedules(enabled_only=True) logger.info(f"Loading {len(schedules)} enabled schedules") for schedule in schedules: try: self._schedule_item(schedule) except Exception as e: logger.error(f"Error scheduling '{schedule['name']}': {e}") def _schedule_item(self, schedule: dict): """ Schedule a job or group with APScheduler Args: schedule: Schedule definition from database """ import json schedule_id = schedule['id'] schedule_name = schedule['name'] target_type = schedule['target_type'] target_id = schedule['target_id'] schedule_type = schedule['schedule_type'] # Parse schedule config try: schedule_config = json.loads(schedule['schedule_config']) except: schedule_config = schedule['schedule_config'] # Skip manual schedules (no automatic execution) if schedule_type == 'manual': logger.debug(f"Skipping manual schedule '{schedule_name}'") return # Build trigger based on schedule type trigger = self._build_trigger(schedule_type, schedule_config) if not trigger: logger.warning(f"Could not build trigger for schedule '{schedule_name}', skipping") return # Determine execution function based on target type if target_type == 'job': exec_func = self._execute_scheduled_job exec_args = [target_id] elif target_type == 'group': exec_func = self._execute_scheduled_group exec_args = [target_id] else: logger.error(f"Unknown target_type '{target_type}' in schedule '{schedule_name}'") return # Add to APScheduler self.scheduler.add_job( func=exec_func, trigger=trigger, args=exec_args, id=f"schedule_{schedule_id}", name=schedule_name, replace_existing=True ) logger.info(f"Scheduled '{schedule_name}' (ID {schedule_id}): {target_type} {target_id} with {schedule_type}") def _build_trigger(self, schedule_type: str, schedule_config: dict): """ Build APScheduler trigger from schedule configuration Args: schedule_type: 'cron', 'interval', or 'date' schedule_config: Configuration dict Returns: APScheduler trigger object """ if schedule_type == 'cron': # Cron trigger: {'cron': '0 7 * * mon-fri'} cron_expr = schedule_config.get('cron') if cron_expr: # Convert Unix day-of-week (0=Sun) to ISO day-of-week (0=Mon) cron_expr_iso = convert_unix_to_iso_dayofweek(cron_expr) return CronTrigger.from_crontab(cron_expr_iso) elif schedule_type == 'interval': # Interval trigger: {'hours': 1} or {'minutes': 30} return IntervalTrigger(**schedule_config) elif schedule_type == 'date': # One-time trigger: {'run_date': '2025-11-26 10:00:00'} run_date = schedule_config.get('run_date') if run_date: return DateTrigger(run_date=run_date) return None def _execute_scheduled_job(self, job_id: int): """ Execute a scheduled job (called by APScheduler) Args: job_id: Job ID to execute """ logger.info(f"Triggered scheduled execution of job ID {job_id}") try: execution_id = self.executor.execute_job(job_id, triggered_by='scheduler') if execution_id: logger.info(f"Job ID {job_id} queued for execution (execution ID {execution_id})") else: logger.warning(f"Job ID {job_id} could not be queued") except Exception as e: logger.exception(f"Error executing scheduled job {job_id}: {e}") def _execute_scheduled_group(self, group_id: int): """ Execute a scheduled group (called by APScheduler) Args: group_id: Group ID to execute """ logger.info(f"Triggered scheduled execution of group ID {group_id}") try: # Use DependencyManager to execute group in sequence result = self.dependency_manager.execute_job_group( group_id, triggered_by='scheduler' ) if result.get('success'): logger.info(f"Group ID {group_id} completed successfully") else: failed_count = len(result.get('failed_jobs', [])) logger.warning(f"Group ID {group_id} completed with {failed_count} failed job(s)") except Exception as e: logger.exception(f"Error executing scheduled group {group_id}: {e}") def _health_check(self): """Periodic health check and cleanup""" # Frank-relay: gestisci richieste sola-lettura da Nave via Dropbox try: process_pending_requests() except Exception as e: logger.warning(f"Frank-relay: errore nel ciclo di polling: {e}") # Check for any stuck jobs active = self.db.get_active_executions() for execution in active: # Check if execution has been running too long start_time = datetime.fromisoformat(execution['start_time']) running_time = (datetime.now() - start_time).total_seconds() # If running more than 2x timeout, mark as failed job = self.db.get_job(execution['job_id']) if job and running_time > (job['timeout_seconds'] * 2): logger.warning(f"Execution {execution['id']} appears stuck, marking as failed") self.db.update_execution( execution['id'], status='failed', end_time=datetime.now(), error_message='Execution exceeded 2x timeout - marked as failed by health check' ) def add_job_to_schedule(self, job_id: int): """ Add or update a job in the scheduler Args: job_id: Job ID to add/update """ job = self.db.get_job(job_id) if job and job['enabled']: self._schedule_job(job) logger.info(f"Added/updated job '{job['name']}' in scheduler") def remove_job_from_schedule(self, job_id: int): """ Remove a job from the scheduler Args: job_id: Job ID to remove """ try: self.scheduler.remove_job(f"job_{job_id}") logger.info(f"Removed job ID {job_id} from scheduler") except Exception as e: logger.warning(f"Could not remove job {job_id} from scheduler: {e}") def execute_job_manually(self, job_id: int) -> Optional[int]: """ Execute a job manually (outside of schedule) Args: job_id: Job ID to execute Returns: Execution ID """ logger.info(f"Manual execution requested for job ID {job_id}") return self.executor.execute_job(job_id, triggered_by='manual') def get_scheduler_status(self) -> dict: """Get scheduler status information""" jobs = self.scheduler.get_jobs() active_executions = self.db.get_active_executions() return { 'running': self.running, 'scheduled_jobs': len(jobs), 'active_executions': len(active_executions), 'next_run_time': str(jobs[0].next_run_time) if jobs else None } def main(): """Main entry point for daemon""" logger.info("="*60) logger.info("Python Scheduler Daemon - Starting") logger.info("="*60) daemon = SchedulerDaemon() try: daemon.start() except KeyboardInterrupt: logger.info("Received KeyboardInterrupt") except Exception as e: logger.exception(f"Fatal error in daemon: {e}") sys.exit(1) sys.exit(0) if __name__ == '__main__': main()