""" Automation Daemon - Unified daemon for Frank Manages scheduling and execution for jobs and backup modules """ import time import signal import socket import sys from typing import Optional, Dict, Any 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 Database 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 AutomationDaemon: """Unified daemon managing multiple automation modules""" def __init__(self): # Setup logging first setup_logging('daemon') logger.info("Initializing Frank Automation Daemon") # Current hostname — used for job filtering self.hostname = socket.gethostname() logger.info(f"Running on host: {self.hostname}") # Load configuration self.config = SchedulerConfig() # Initialize database — split mode: local runtime DB + remote B for config sync self.db = Database(self.config.database_path, self.config.local_database_path) # Initialize notification manager self.notification_manager = NotificationManager(self.db, self.config) # Initialize scheduler module components self.executor = JobExecutor(self.db, self.config, self.notification_manager) self.dependency_manager = DependencyManager(self.db, self.executor) # Initialize modules registry self.modules: Dict[str, Any] = {} self._load_modules() # Initialize APScheduler self.scheduler = BackgroundScheduler( timezone='Europe/Rome', job_defaults={ 'coalesce': True, # Combine multiple missed runs into one 'max_instances': 1, # Prevent overlapping executions 'misfire_grace_time': 3600 # 1 hour — covers machine sleep/wake around scheduled time } ) # Shutdown flag self.running = False # Stealth mode flag (outside active hours) self._stealth_mode = False # Config version polling (remote B) self._remote_config_version: Optional[int] = None self._last_version_check: float = 0.0 self._version_check_interval: int = 60 # seconds # Setup signal handlers signal.signal(signal.SIGINT, self._signal_handler) signal.signal(signal.SIGTERM, self._signal_handler) logger.info("Automation Daemon initialized") def _load_modules(self): """Load and initialize automation modules""" logger.info("Loading automation modules...") # Scheduler module (always available, built-in) self.modules['scheduler'] = { 'name': 'Scheduler', 'executor': self.executor, 'dependency_manager': self.dependency_manager } logger.info(" ✓ Scheduler module loaded") # Backup module (optional, load if available) try: from modules.backup import BackupModule self.modules['backup'] = { 'name': 'Backup', 'instance': BackupModule(self.db) } logger.info(" ✓ Backup module loaded") except Exception as e: logger.warning(f" ✗ Backup module not available: {e}") # Future modules can be loaded here # try: # from modules.monitoring import MonitoringModule # self.modules['monitoring'] = MonitoringModule(self.db) # except: # pass logger.info(f"Loaded {len(self.modules)} module(s)") def _signal_handler(self, signum, frame): """Handle shutdown signals""" logger.info(f"Received signal {signum}, shutting down...") self.shutdown() def start(self): """Start the automation daemon""" logger.info("Starting Automation Daemon") try: # Load all enabled schedules from database self._load_schedules() # Start APScheduler self.scheduler.start() self.running = True logger.info("Automation Daemon started successfully") logger.info(f"Loaded {len(self.scheduler.get_jobs())} schedule(s)") # Main loop - keep daemon alive try: while self.running: time.sleep(self.config.get('daemon.check_interval_seconds', 30)) # Frank-relay: gestisci richieste sola-lettura da Nave (sempre attivo, anche in stealth) try: process_pending_requests() except Exception as e: logger.warning(f"Frank-relay: errore nel ciclo di polling: {e}") if not self._is_active_window(): if not self._stealth_mode: self._enter_stealth() continue if self._stealth_mode: self._exit_stealth() # Periodic health check self._health_check() # Poll remote B for config changes (every 60s) now = time.monotonic() if now - self._last_version_check >= self._version_check_interval: self._check_remote_config_version() self._last_version_check = now 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 Automation 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") # Cleanup modules for module_name, module_data in self.modules.items(): if 'instance' in module_data and hasattr(module_data['instance'], 'cleanup'): try: module_data['instance'].cleanup() logger.info(f"Module '{module_name}' cleanup completed") except Exception as e: logger.error(f"Error cleaning up module '{module_name}': {e}") logger.info("Automation 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 schedule(s)") 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 an item 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': # Check if this job belongs to the current host job = self.db.get_job(target_id) if job and job.get('hostname', '') and job['hostname'] != self.hostname: logger.debug(f"Skipping schedule '{schedule_name}': job hostname='{job['hostname']}', current='{self.hostname}'") return exec_func = self._execute_scheduled_job exec_args = [target_id] elif target_type == 'group': # Check that all jobs in the group belong to this host group_jobs = self.db.get_job_group_members(target_id) for gj in group_jobs: job_hostname = gj.get('hostname', '') if job_hostname and job_hostname != self.hostname: logger.debug(f"Skipping group schedule '{schedule_name}': job '{gj['job_name']}' hostname='{job_hostname}', current='{self.hostname}'") return exec_func = self._execute_scheduled_group exec_args = [target_id] elif target_type == 'backup_profile': # Backup module: execute backup if 'backup' not in self.modules: logger.error(f"Cannot schedule backup: backup module not loaded") return exec_func = self._execute_scheduled_backup exec_args = [schedule_id, 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 _execute_scheduled_backup(self, schedule_id: int, profile_id: int): """ Execute a scheduled backup (called by APScheduler) Args: schedule_id: Schedule ID that triggered profile_id: Backup profile ID to execute """ logger.info(f"Triggered scheduled execution of backup profile ID {profile_id}") try: backup_module = self.modules['backup']['instance'] result = backup_module.schedule_item(schedule_id, profile_id) if result.get('success'): logger.info(f"Backup profile ID {profile_id} completed successfully") else: error = result.get('error', 'Unknown error') logger.error(f"Backup profile ID {profile_id} failed: {error}") except Exception as e: logger.exception(f"Error executing scheduled backup {profile_id}: {e}") def _is_active_window(self) -> bool: """Returns True if current time is within active window: Mon-Fri 8-20""" now = datetime.now() if now.weekday() >= 5: # 5=Saturday, 6=Sunday return False start = self.config.get('daemon.active_hours_start', 8) end = self.config.get('daemon.active_hours_end', 20) return start <= now.hour < end def _enter_stealth(self): """Pause all activity outside active hours — no DB access, no job execution""" self._stealth_mode = True if self.scheduler.running: self.scheduler.pause() start = self.config.get('daemon.active_hours_start', 8) logger.info(f"Stealth mode ON — nessuna attività fino alle {start:02d}:00 (ora corrente: {datetime.now().strftime('%H:%M')})") def _exit_stealth(self): """Resume normal activity at start of active hours""" self._stealth_mode = False if self.scheduler.running: self.scheduler.resume() logger.info(f"Stealth mode OFF — ripresa attività (ora: {datetime.now().strftime('%H:%M')})") def _health_check(self): """Periodic health check and cleanup""" # 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 _check_remote_config_version(self): """ Poll remote B for config_version changes. If version changed → re-sync config tables + reload schedules (no restart needed). If B unreachable → skip silently (no crash). """ v = self.db.get_remote_config_version() if v is None: return # B unreachable — ignore, retry next cycle if self._remote_config_version is not None and v != self._remote_config_version: logger.info(f"Config version changed on remote B ({self._remote_config_version} → {v}) — resyncing") synced = self.db._sync_config_from_remote() if synced: self._reload_schedules() self._remote_config_version = v def _reload_schedules(self): """Remove all APScheduler jobs and reload from local DB (after config sync).""" logger.info("Reloading schedules from local DB") for job in self.scheduler.get_jobs(): job.remove() self._load_schedules() logger.info(f"Schedules reloaded: {len(self.scheduler.get_jobs())} active") def add_schedule(self, schedule_id: int): """ Add or update a schedule in the scheduler Args: schedule_id: Schedule ID to add/update """ schedule = self.db.get_schedule(schedule_id) if schedule and schedule['enabled']: self._schedule_item(schedule) logger.info(f"Added/updated schedule '{schedule['name']}' in scheduler") def remove_schedule(self, schedule_id: int): """ Remove a schedule from the scheduler Args: schedule_id: Schedule ID to remove """ try: self.scheduler.remove_job(f"schedule_{schedule_id}") logger.info(f"Removed schedule ID {schedule_id} from scheduler") except Exception as e: logger.warning(f"Could not remove schedule {schedule_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 execute_backup_manually(self, profile_id: int) -> Dict[str, Any]: """ Execute a backup manually (outside of schedule) Args: profile_id: Backup profile ID to execute Returns: Result dict with success status and details """ logger.info(f"Manual execution requested for backup profile ID {profile_id}") if 'backup' not in self.modules: error_msg = "Backup module not loaded" logger.error(error_msg) return {"success": False, "error": error_msg} backup_module = self.modules['backup']['instance'] return backup_module.execute_backup(profile_id) def get_daemon_status(self) -> dict: """Get daemon status information""" jobs = self.scheduler.get_jobs() active_executions = self.db.get_active_executions() # Get module-specific status modules_status = {} for module_name, module_data in self.modules.items(): modules_status[module_name] = { 'loaded': True, 'name': module_data.get('name', module_name) } return { 'running': self.running, 'scheduled_jobs': len(jobs), 'active_executions': len(active_executions), 'modules': modules_status, '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 Automation Daemon - Starting") logger.info("="*60) daemon = AutomationDaemon() 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()