"""Состояние станции: ``GET /healthz`` и его сборка.
Процесс с мёртвым планировщиком снаружи неотличим от здорового — эндпоинт
закрывает именно это. Отдаёт состояние планировщика, время последней успешной
сверки расписания, глубину очереди ``incomplete/``, признак идущего прохода и
время следующего.
Ходят сюда двое: ``HEALTHCHECK`` образа и агент управления парком из соседнего
контейнера.
"""
import json
import threading
from collections.abc import Callable
from datetime import UTC, datetime
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from apscheduler.schedulers.base import BaseScheduler
from core.configs import logger, settings
from soniks_client.scheduler import SYNC_JOB_ID, find_next_observation
# Состояние модульное и без локов — так и задумано: присваивание и
# ``add``/``discard`` атомарны под GIL, а наружу отдаётся только ``bool``, без
# обхода множества (обход из потока HTTP параллельно с ``add`` из пула
# планировщика дал бы RuntimeError).
_last_sync: datetime | None = None
_active_observations: set[str] = set()
_clock_skew: int | None = None
_config_generation: int = 0
_get_scheduler: Callable[[], BaseScheduler | None] | None = None
[документация]
def mark_sync_success() -> None:
"""Отметить успешно завершённый цикл сверки расписания."""
global _last_sync
_last_sync = datetime.now(UTC)
[документация]
def set_clock_skew(seconds: int | None) -> None:
"""Запомнить расхождение часов станции с порталом.
Плюс — часы станции спешат. ``None`` — портал не прислал ``Date``.
"""
global _clock_skew
_clock_skew = seconds
[документация]
def set_config_generation(generation: int) -> None:
"""Запомнить поколение конфигурации с портала, применённое станцией."""
global _config_generation
_config_generation = generation
[документация]
def observation_running() -> bool:
"""Идёт ли сейчас хотя бы один проход."""
return bool(_active_observations)
[документация]
def seconds_to_next_observation() -> float | None:
"""Сколько секунд до ближайшего прохода; ``None`` — проходов нет.
По этому окну обследование приёмника решает, можно ли его занять
(``soniks_client.sdr_survey``). Идущий проход — ноль.
"""
if _active_observations:
return 0.0
scheduler = _get_scheduler() if _get_scheduler is not None else None
if scheduler is None or not scheduler.running:
return None
job = find_next_observation(scheduler)
if job is None:
return None
return (job.kwargs["start"] - datetime.now(UTC)).total_seconds()
[документация]
def observation_started(job_id: str) -> None:
"""Отметить начало прохода.
Множество, а не флаг и не мьютекс: наложившиеся проходы разрешены
осознанно (лок на них отклонён в дорожной карте техдолга), и вложенный
проход не должен гасить признак у соседнего.
"""
_active_observations.add(job_id)
[документация]
def observation_finished(job_id: str) -> None:
"""Отметить завершение прохода."""
_active_observations.discard(job_id)
def _incomplete_count() -> int:
"""Посчитать наблюдения, ожидающие повторной выгрузки."""
try:
return sum(
1 for path in settings.paths.incomplete_path.iterdir() if path.is_dir()
)
except OSError:
return 0
[документация]
def snapshot(
get_scheduler: Callable[[], BaseScheduler | None],
) -> tuple[dict, int]:
"""Собрать состояние станции и код ответа.
Returns:
Тело ответа и код HTTP: 503, если планировщик не запущен или потеряно
задание сверки расписания.
"""
scheduler = get_scheduler()
scheduler_running = scheduler is not None and scheduler.running
sync_job_alive = False
next_observation = None
if scheduler_running:
sync_job_alive = scheduler.get_job(SYNC_JOB_ID) is not None
next_observation_job = find_next_observation(scheduler)
if next_observation_job is not None:
next_observation = next_observation_job.kwargs["start"].isoformat()
last_sync = _last_sync
last_sync_age = None
if last_sync is not None:
last_sync_age = int((datetime.now(UTC) - last_sync).total_seconds())
body = {
"scheduler_running": scheduler_running,
"last_sync": last_sync.isoformat() if last_sync is not None else None,
"last_sync_age_seconds": last_sync_age,
"incomplete": _incomplete_count(),
"observation_running": bool(_active_observations),
"next_observation": next_observation,
"clock_skew_seconds": _clock_skew,
"config_generation": _config_generation,
}
# Неудачная сверка на код ответа не влияет: лежащий портал — не повод
# считать станцию больной, чинить её перезапуском нечего. Возраст
# ``last_sync`` и уход часов отданы как данные, решение принимает
# потребитель.
status = 200 if scheduler_running and sync_job_alive else 503
return body, status
class _HealthHandler(BaseHTTPRequestHandler):
"""Единственный маршрут ``GET /healthz``."""
def do_GET(self) -> None:
if self.path.split("?")[0] != "/healthz":
self._respond({"detail": "not found"}, 404)
return
try:
body, status = snapshot(self.server.get_scheduler)
except Exception as e:
# Без перехвата исключение уходит в handle_error: соединение
# рвётся без ответа, и проверяющий получает таймаут вместо кода.
logger.exception("Ошибка сборки состояния для /healthz: %s", e)
body, status = {"detail": str(e)}, 503
self._respond(body, status)
def _respond(self, body: dict, status: int) -> None:
payload = json.dumps(body).encode()
self.send_response(status)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(payload)))
self.end_headers()
self.wfile.write(payload)
def log_message(self, format: str, *args) -> None:
# Иначе каждый запрос уходит в stderr мимо логгера приложения.
logger.debug("healthz: " + format, *args)
[документация]
def start_health_server(
get_scheduler: Callable[[], BaseScheduler | None],
) -> ThreadingHTTPServer | None:
"""Поднять ``/healthz`` в фоновом потоке.
Args:
get_scheduler: Колбэк, а не сам планировщик: супервизор подменяет
экземпляр при перезапуске, и захваченная ссылка устарела бы.
Returns:
Запущенный сервер или ``None``, если порт занять не удалось.
"""
global _get_scheduler
# Планировщик нужен не только /healthz: по нему же считается окно до
# ближайшего прохода для обследования приёмника.
_get_scheduler = get_scheduler
try:
# 0.0.0.0, а не 127.0.0.1: состояние читает агент из соседнего
# контейнера. Наружу хоста не торчит, пока в compose нет ни ``ports:``,
# ни ``network_mode: host`` — появится любое из двух, переводить на
# 127.0.0.1.
server = ThreadingHTTPServer(
("0.0.0.0", settings.health.PORT),
_HealthHandler,
)
except OSError as e:
# Занятый порт не повод не проводить проходы.
logger.error(
"Не удалось поднять /healthz на порту %d: %s", settings.health.PORT, e
)
return None
server.get_scheduler = get_scheduler
threading.Thread(target=server.serve_forever, name="healthz", daemon=True).start()
logger.info("Эндпоинт /healthz слушает порт %d", server.server_port)
return server