"""Отдельный процесс-сборщик: держит расписание и запускает сбор ASN/FQDN. Общение с API - через файлы: расписание читается из config.json (его меняет POST /schedule), состояние заданий и heartbeat пишутся в status.json (его читает GET /health). """ import datetime import logging import os import signal import sys import threading from apscheduler.schedulers.blocking import BlockingScheduler from apscheduler.triggers.cron import CronTrigger import cidr_collector as cc import db from cidr_collector import CIDRCollector, FQDNCollector, load_full_config from storage import StorageError, save_json_atomic, try_lock logger = logging.getLogger("collector_daemon") DEFAULT_CRONS = {"asn": "0 2 * * *", "fqdn": "0 3 * * *", "backup": "30 4 * * *"} def run_backup(): """Задание backup: онлайн-копия базы, проверка, ротация (число копий - backup_keep в config.json).""" keep = max(1, int(load_full_config().get("backup_keep", cc.DEFAULT_BACKUP_KEEP))) with db.session() as conn: path = db.backup_database(conn, datetime.datetime.now(), keep) if path is None: logger.warning("Database backup skipped: no data yet although sources are configured") return logger.info("Database backup written to %s (keeping the last %d)", path, keep) # Задания: имя -> вызываемый объект (новый экземпляр на каждый запуск - свежий config.json) COLLECTORS = {"asn": lambda: CIDRCollector().run_collection(), "fqdn": lambda: FQDNCollector().run_collection(), "backup": run_backup} LOCK_FILE = os.path.join(cc.DATA_DIR, "collector.daemon.lock") # Состояние заданий: действующий cron, результат последнего запуска, отклонённый cron (чтобы не спамить логом) job_state = {name: {"cron": None, "last_run": None, "last_finished": None, "running": False, "last_error": None, "rejected_cron": None} for name in COLLECTORS} # RLock: изменения состояния и write_status (читает под той же блокировкой) не должны мешать друг другу _state_lock = threading.RLock() # Один запуск за раз на тип: наложение планового и ручного сбора пропускается _run_locks = {name: threading.Lock() for name in COLLECTORS} def write_status(scheduler): """Атомарно пишет status.json; вызывается после каждого запуска и как heartbeat.""" with _state_lock: jobs = {} for name, state in job_state.items(): job = scheduler.get_job(f"{name}_job") next_run = getattr(job, "next_run_time", None) if job else None jobs[name] = {"cron": state["cron"], "last_run": state["last_run"], "last_finished": state["last_finished"], "running": state["running"], "last_error": state["last_error"], "next_run": next_run.isoformat() if next_run else None} save_json_atomic(cc.STATUS_FILE, {"updated_at": datetime.datetime.now().isoformat(), "jobs": jobs}) def run_job(name, scheduler): lock = _run_locks[name] if not lock.acquire(blocking=False): logger.warning("%s collection is already running, skipping this start", name) return try: logger.info("Running %s collection...", name) state = job_state[name] with _state_lock: error = state["last_error"] # при прерывании (SystemExit) прежняя ошибка остаётся state["last_run"] = datetime.datetime.now().isoformat() state["running"] = True write_status(scheduler) try: COLLECTORS[name]() error = None except Exception as e: logger.exception("%s collection failed", name) error = str(e) finally: with _state_lock: state["last_error"] = error state["running"] = False state["last_finished"] = datetime.datetime.now().isoformat() write_status(scheduler) finally: lock.release() def check_collect_requests(scheduler): """Забирает запрос POST /collect и запускает сбор немедленно (разовыми заданиями).""" try: types = cc.pop_collection_requests() except StorageError as e: logger.error("Cannot read collect request: %s", e) return for name in types: logger.info("Manual %s collection requested", name) scheduler.add_job(run_job, args=[name, scheduler], id=f"{name}_manual", replace_existing=True) def schedule_jobs(scheduler): """Создаёт задания по расписанию из config.json (для отсутствующих значений - по умолчанию).""" schedule = load_full_config().get("schedule", {}) for name in COLLECTORS: cron = schedule.get(name, DEFAULT_CRONS[name]) scheduler.add_job(run_job, CronTrigger.from_crontab(cron), args=[name, scheduler], id=f"{name}_job", replace_existing=True, max_instances=1, coalesce=True) with _state_lock: job_state[name]["cron"] = cron logger.info("Job %s scheduled: %s", name, cron) def sync_schedule(scheduler): """Подхватывает изменения расписания из config.json без перезапуска; пишет heartbeat.""" try: schedule = load_full_config().get("schedule", {}) except StorageError as e: logger.error("Cannot read config, keeping current schedule: %s", e) schedule = {} for name in COLLECTORS: cron = schedule.get(name) state = job_state[name] if cron is None or cron == state["cron"] or cron == state["rejected_cron"]: continue try: trigger = CronTrigger.from_crontab(cron) except ValueError as e: logger.error("Invalid cron for %s (%r), keeping %r: %s", name, cron, state["cron"], e) with _state_lock: state["rejected_cron"] = cron continue scheduler.reschedule_job(f"{name}_job", trigger=trigger) logger.info("Job %s rescheduled: %s -> %s", name, state["cron"], cron) with _state_lock: state["cron"] = cron state["rejected_cron"] = None write_status(scheduler) def build_scheduler(): scheduler = BlockingScheduler() schedule_jobs(scheduler) scheduler.add_job(sync_schedule, "interval", seconds=cc.SYNC_INTERVAL, args=[scheduler], id="sync_schedule", max_instances=1, coalesce=True) scheduler.add_job(check_collect_requests, "interval", seconds=cc.TRIGGER_POLL_INTERVAL, args=[scheduler], id="collect_requests", max_instances=1, coalesce=True) return scheduler def main(): logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s") # Служебные задания (опрос запросов каждые 5 с) не должны засорять лог logging.getLogger("apscheduler.executors.default").setLevel(logging.WARNING) lock = try_lock(LOCK_FILE) if lock is None: logger.error("Another collector daemon is already running (%s is locked).", LOCK_FILE) sys.exit(1) try: scheduler = build_scheduler() except (StorageError, ValueError) as e: logger.error("Cannot start: %s", e) sys.exit(1) # SIGTERM -> штатный выход: текущий сбор завершится, файлы пишутся атомарно signal.signal(signal.SIGTERM, lambda signum, frame: sys.exit(0)) write_status(scheduler) logger.info("Collector daemon started.") try: scheduler.start() except (KeyboardInterrupt, SystemExit): pass finally: if scheduler.running: scheduler.shutdown(wait=False) logger.info("Collector daemon stopped.") if __name__ == "__main__": main()