Files
ripe-cidr-collector/collector_daemon.py
T
ayurishchevandClaude Sonnet 5 aaf41adc79 Stage B of review fixes: daemon state lock, image build check, RIPEstat retries
Job state is changed under one RLock, the Dockerfile copies all root
modules and imports them at build time, RIPEstat requests go through a
retrying session with the sourceapp parameter (ripestat_sourceapp) and a
capped Retry-After. Adds the summary and marks review findings 5-10 fixed.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-21 10:12:48 +03:00

189 lines
8.1 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}
# 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()