Исходный код soniks_client.scheduler
"""Фабрика ``BackgroundScheduler`` (UTC), собранного из настроек."""
import logging
from apscheduler.executors.pool import ThreadPoolExecutor
from apscheduler.job import Job
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.schedulers.base import BaseScheduler
from core.configs import settings
logging.getLogger("apscheduler").setLevel(settings.scheduler.LOG_LEVEL)
# Идентификатор постоянного задания сверки расписания. Лежит здесь, а не рядом
# с регистрацией в ``jobs/__init__.py``: его читает ещё и ``/healthz``, а импорт
# оттуда из пакета заданий замкнул бы кольцо импортов.
SYNC_JOB_ID = "register_observation_jobs"
[документация]
def create_scheduler() -> BackgroundScheduler:
"""Собрать новый планировщик.
Именно фабрика, а не общий экземпляр: APScheduler после ``shutdown()``
навсегда гасит пул потоков исполнителя, поэтому перезапуск того же объекта
возвращает планировщик, который принимает задания и не выполняет их.
"""
return BackgroundScheduler(
executors={"default": ThreadPoolExecutor(settings.scheduler.MAX_WORKERS)},
job_defaults={
"coalesce": settings.scheduler.COALESCE,
"max_instances": settings.scheduler.MAX_INSTANCES,
"misfire_grace_time": settings.scheduler.MISFIRE_GRACE_TIME_IN_SECOND,
},
timezone="UTC",
)
[документация]
def find_next_observation(scheduler: BaseScheduler) -> Job | None:
"""Найти ближайшее запланированное наблюдение.
Returns:
Задание наблюдения с самым ранним началом или ``None``, если проходов
не запланировано.
"""
observation_jobs: list[Job] = [
job
for job in scheduler.get_jobs()
if settings.jobs.is_observation_job(job.name)
]
if not observation_jobs:
return None
# Время берётся из kwargs, которыми задание и создавалось: trigger.run_date
# есть только у DateTrigger, а next_run_time у задания, добавленного в
# остановленный планировщик, отсутствует как атрибут вовсе.
return min(observation_jobs, key=lambda job: job.kwargs["start"])