"""Сверка расписания портала с состоянием планировщика.
Идентификатор задания на портале совпадает с идентификатором задания в
планировщике — на этом строится добавление, обновление и удаление."""
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,
)