Исходный код soniks_client.jobs.sync

"""Сверка расписания портала с состоянием планировщика.

Идентификатор задания на портале совпадает с идентификатором задания в
планировщике — на этом строится добавление, обновление и удаление."""

from datetime import UTC, datetime

from apscheduler.job import Job
from apscheduler.jobstores.base import JobLookupError
from apscheduler.schedulers.base import BaseScheduler

from core.configs import logger, settings
from soniks_client import remote_config, sdr_survey
from soniks_client.api import get_observation_jobs, publish_station_status
from soniks_client.health import mark_sync_success
from soniks_client.jobs.observation import execute_observation
from soniks_client.models import JobData


[документация] def sync_observation_jobs(scheduler: BaseScheduler) -> None: """Свести расписание портала с состоянием планировщика. Задания делятся на добавляемые, обновляемые и удаляемые; отменённый на портале проход исчезает и у станции. """ jobs = get_observation_jobs() # Именно `is None`: пустой список — это «портал снял все проходы», и # синхронизацию надо провести, иначе отменённый последний проход # навсегда остаётся на станции. if jobs is None: return jobs_to_remove, jobs_to_update, jobs_to_add = _divide_tasks_into_categories( jobs, scheduler, ) _synchronize_jobs( jobs_to_remove, jobs_to_update, jobs_to_add, scheduler, ) # Отметка только здесь: ранний выход выше — это недоступный портал, то есть # успешного цикла не было. Пустой список от портала успехом считается. mark_sync_success() # Конфигурация с портала — до публикации статуса: статус уезжает уже с # итогом применения. Отложенное из-за прохода применяется здесь же. remote_config.sync() # Обследование приёмника — в фоне и только в окне без проходов; итог # опубликует сам поток обследования. sdr_survey.run_due() # Здесь, а не на старте: портал только что ответил, а если публикация не # удалась, следующая сверка повторит её сама. publish_station_status()
def _divide_tasks_into_categories( jobs: list[JobData], scheduler: BaseScheduler, ) -> tuple[set[str], dict[str, JobData], dict[str, JobData]]: latest_network_jobs = {job.id: job for job in jobs} all_jobs: list[Job] = scheduler.get_jobs() current_observation_jobs: dict[str, JobData] = { job.id: _create_job_data(job) for job in all_jobs if settings.jobs.is_observation_job(job.name) } latest_job_ids = set(latest_network_jobs.keys()) current_observation_job_ids = set(current_observation_jobs.keys()) new_jobs = latest_job_ids - current_observation_job_ids jobs_to_remove = current_observation_job_ids - latest_job_ids potential_jobs_to_update = latest_job_ids & current_observation_job_ids jobs_to_add: dict[str, JobData] = { job_id: latest_network_jobs[job_id] for job_id in new_jobs } jobs_to_update: dict[str, JobData] = { job_id: latest_network_jobs[job_id] for job_id in potential_jobs_to_update if latest_network_jobs[job_id] != current_observation_jobs[job_id] } logger.debug( "Количество заданий для удаления: %d, для обновления: %d, для добавления: %d", len(jobs_to_remove), len(jobs_to_update), len(jobs_to_add), ) return jobs_to_remove, jobs_to_update, jobs_to_add def _create_job_data(job: Job) -> JobData: """Восстановить ``JobData`` из задания планировщика для сверки с порталом. Все поля берутся из ``kwargs``, которыми задание и создавалось: без ``norad_cat_id`` сравнение с данными портала **никогда** не совпадало и каждый проход удалялся и добавлялся заново раз в минуту. ``start`` тоже берётся из ``kwargs``, а не из ``trigger.run_date`` — это снимает завязку на конкретный тип триггера. """ return JobData( id=job.id, start=job.kwargs["start"], end=job.kwargs["end"], tle=job.kwargs["tle"], frequency=job.kwargs["frequency"], mode=job.kwargs["mode"], baud=job.kwargs["baud"], norad_cat_id=job.kwargs["norad_cat_id"], ) def _synchronize_jobs( jobs_to_remove: set[str], jobs_to_update: dict[str, JobData], jobs_to_add: dict[str, JobData], scheduler: BaseScheduler, ) -> None: _remove_jobs(jobs_to_remove, scheduler) _update_jobs(jobs_to_update, scheduler) _add_jobs(jobs_to_add, scheduler) def _remove_jobs( jobs_to_remove: set[str], scheduler: BaseScheduler, ) -> None: if not jobs_to_remove: return for job_id in jobs_to_remove: if _remove_job(job_id, scheduler): logger.info("Удалено задание [ID: %s]", job_id) def _update_jobs( jobs_to_update: dict[str, JobData], scheduler: BaseScheduler, ) -> None: if not jobs_to_update: return for job_id, job in jobs_to_update.items(): if job.start < datetime.now(tz=UTC): continue if not _remove_job(job_id, scheduler): continue _add_job_to_scheduler(job, scheduler) logger.info("Обновлено задание [ID: %s]: %s", job_id, job) def _remove_job(job_id: str, scheduler: BaseScheduler) -> bool: """Снять задание с планировщика, пережив уже отработавшее. ``remove_job`` кидает ``JobLookupError``, если задание только что отработало и само себя сняло; раньше это исключение обрывало остаток цикла синхронизации. """ try: scheduler.remove_job(job_id) except JobLookupError: logger.debug("Задание [ID: %s] уже снято с планировщика", job_id) return False return True def _add_jobs( jobs_to_add: dict[str, JobData], scheduler: BaseScheduler, ) -> None: if not jobs_to_add: return for job_id, job in jobs_to_add.items(): _add_job_to_scheduler(job, scheduler) logger.info("Добавлено новое задание [ID: %s]: %s", job_id, job) def _add_job_to_scheduler( job: JobData, scheduler: BaseScheduler, ) -> None: scheduler.add_job( func=execute_observation, trigger="date", run_date=job.start, id=job.id, name=settings.jobs.get_job_name(job.id), kwargs={ "job_id": job.id, "tle": job.tle, "start": job.start, "end": job.end, "frequency": job.frequency, "mode": job.mode, "baud": job.baud, "norad_cat_id": job.norad_cat_id, "scheduler": scheduler, }, replace_existing=True, )