Daemon job "backup" (online copy, quick_check, rotation), restore of the newest valid copy when the database cannot be opened, change journal reset after restore, last_restore in /health, docs and tests. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
181 lines
7.7 KiB
Python
181 lines
7.7 KiB
Python
"""Отдельный процесс-сборщик: держит расписание и запускает сбор 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)
|
|
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}
|
|
_state_lock = threading.Lock()
|
|
# Один запуск за раз на тип: наложение планового и ручного сбора пропускается
|
|
_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]
|
|
state["last_run"] = datetime.datetime.now().isoformat()
|
|
state["running"] = True
|
|
write_status(scheduler)
|
|
try:
|
|
COLLECTORS[name]()
|
|
state["last_error"] = None
|
|
except Exception as e:
|
|
logger.exception("%s collection failed", name)
|
|
state["last_error"] = str(e)
|
|
finally:
|
|
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)
|
|
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)
|
|
state["rejected_cron"] = cron
|
|
continue
|
|
scheduler.reschedule_job(f"{name}_job", trigger=trigger)
|
|
logger.info("Job %s rescheduled: %s -> %s", name, state["cron"], cron)
|
|
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()
|