Исходный код 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"])