Files
KindClustersDashboard/app/core/job_store.py
T
Sergey Antropoff 4546f50aef UI: навигация, документация, favicon; журнал развёртывания; валидация формы
- Меню API: ссылка «Теги образов», обёртка прокрутки пилюль, z-index и padding против обрезки hover
- Документация: ширина колонки как у дашборда (72rem)
- Favicon SVG + GET /favicon.ico, link в base.html
- provision_log.json, GET .../provision-log, кнопка в таблице кластеров
- Валидация create: имя, workers, тег kindest; модалка alert
- Прочие правки из сессии (clusters, job_store, стили, шаблоны)
2026-04-04 08:15:15 +03:00

455 lines
18 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Хранилище фоновых заданий (создание/старт/стоп кластера).
Снимок заданий сохраняется в JSON в каталоге ``clusters/`` (том на хосте), чтобы после
перезапуска контейнера история не терялась. Записи ``queued``/``running`` при старте
помечаются как ``failed`` — рабочий процесс уже не существует.
Потокобезопасные флаги отмены и прогресс (для worker-thread) — через ``threading.Lock``.
Автор: Сергей Антропов
Сайт: https://devops.org.ru
"""
from __future__ import annotations
import asyncio
import json
import logging
import os
import signal
import subprocess
import threading
import uuid
from collections import deque
from dataclasses import dataclass, field
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Literal, cast
from kind_k8s_paths import clusters_dir
logger = logging.getLogger("kind_k8s.job_store")
# Лимит записей в памяти и в файле (dev-инструмент; старые задания вытесняются)
_MAX_JOBS = 200
# Версия формата JSON на диске
_JOBS_FILE_VERSION = 1
JobStatus = Literal["queued", "running", "success", "failed", "cancelled"]
# --- Синхронное сопровождение задания (worker-thread и HTTP отмена) ---
_thread_lock = threading.Lock()
_cancel_events: dict[str, threading.Event] = {}
_progress: dict[str, tuple[str, int]] = {}
# Хвост логов для активных заданий (kind create и т.д.); после завершения копируется в JobRecord.log_lines
_job_log_deques: dict[str, deque[str]] = {}
# Полный журнал без лимита строк (только для зарегистрированных job_id) — пишется в clusters/<имя>/provision_log.json
_job_uncapped_logs: dict[str, list[str]] = {}
# Активный дочерний процесс (pull / kind create) — для принудительной отмены
_process_lock = threading.Lock()
_active_subprocess_by_job: dict[str, subprocess.Popen] = {}
def _jobs_file_path() -> Path:
"""
Путь к JSON с историей заданий.
Переопределение: ``KIND_K8S_JOBS_JSON`` (абсолютный или относительный путь).
По умолчанию: ``<clusters_dir>/kind_k8s_jobs.json``.
"""
override = (os.environ.get("KIND_K8S_JOBS_JSON") or "").strip()
if override:
return Path(override).expanduser().resolve()
return clusters_dir() / "kind_k8s_jobs.json"
def _max_job_log_lines() -> int:
# Буфер в памяти: старые строки вытесняются; по умолчанию 2500 — виден почти весь типичный лог.
raw = (os.environ.get("KIND_K8S_JOB_LOG_MAX_LINES") or "2500").strip()
try:
return max(50, min(int(raw), 5000))
except ValueError:
return 500
def register_uncapped_job_log(job_id: str) -> None:
"""Включить накопление полного журнала по ``job_id`` (создание/старт кластера → JSON в каталоге кластера)."""
with _thread_lock:
_job_uncapped_logs[job_id] = []
def take_uncapped_log_finalize(job_id: str) -> list[str]:
"""Забрать полный журнал и снять регистрацию (вызывать один раз при завершении задания)."""
with _thread_lock:
lst = _job_uncapped_logs.pop(job_id, None)
return list(lst) if lst else []
def append_log_sync(job_id: str, line: str) -> None:
"""Добавить строку в журнал задания (вызывается из worker-thread во время долгих команд)."""
text = (line or "").rstrip()
if not text:
return
cap = _max_job_log_lines()
with _thread_lock:
uncapped = _job_uncapped_logs.get(job_id)
if uncapped is not None:
uncapped.append(text)
if job_id not in _job_log_deques:
_job_log_deques[job_id] = deque(maxlen=cap)
_job_log_deques[job_id].append(text)
def get_logs_snapshot_sync(job_id: str) -> list[str]:
"""Снимок текущего журнала (для API во время running/queued)."""
with _thread_lock:
d = _job_log_deques.get(job_id)
return list(d) if d else []
def take_logs_finalize_sync(job_id: str) -> list[str]:
"""
Забрать журнал в список и удалить deque (после успеха/ошибки/отмены).
Вызывать перед или внутри обновления JobRecord.
"""
with _thread_lock:
d = _job_log_deques.pop(job_id, None)
return list(d) if d else []
def begin_job_tracking(job_id: str) -> None:
"""Зарегистрировать отмену/прогресс для нового job_id (вызывать при создании задания)."""
with _thread_lock:
_cancel_events[job_id] = threading.Event()
_progress[job_id] = ("В очереди", 0)
def end_job_tracking(job_id: str) -> None:
"""Очистить служебные структуры после завершения задания."""
with _process_lock:
_active_subprocess_by_job.pop(job_id, None)
with _thread_lock:
_cancel_events.pop(job_id, None)
_progress.pop(job_id, None)
_job_log_deques.pop(job_id, None)
_job_uncapped_logs.pop(job_id, None)
def set_progress_sync(job_id: str, stage: str, percent: int) -> None:
"""Обновить текст этапа и процент (0–100) из worker-thread."""
pct = max(0, min(100, int(percent)))
with _thread_lock:
if job_id in _progress:
_progress[job_id] = (stage, pct)
def get_progress_sync(job_id: str) -> tuple[str, int] | None:
"""Снимок прогресса для ответа API."""
with _thread_lock:
return _progress.get(job_id)
def register_subprocess_sync(job_id: str, proc: subprocess.Popen) -> None:
"""Привязать текущий дочерний процесс к заданию (pull, kind create)."""
with _process_lock:
_active_subprocess_by_job[job_id] = proc
def unregister_subprocess_sync(job_id: str) -> None:
"""Снять регистрацию процесса (после завершения потока чтения)."""
with _process_lock:
_active_subprocess_by_job.pop(job_id, None)
def kill_subprocess_for_job_sync(job_id: str) -> None:
"""Принудительно завершить дочерний процесс задания (SIGKILL / группа процессов)."""
with _process_lock:
proc = _active_subprocess_by_job.get(job_id)
if proc is None:
return
if proc.poll() is not None:
return
pid = proc.pid
try:
try:
os.killpg(os.getpgid(pid), signal.SIGKILL)
logger.warning("Отмена: SIGKILL группы процесса PID %s (задание %s)", pid, job_id)
except (AttributeError, OSError, ProcessLookupError):
proc.kill()
logger.warning("Отмена: kill PID %s (задание %s)", pid, job_id)
except OSError as e:
logger.warning("Не удалось убить процесс задания %s: %s", job_id, e)
def request_cancel_sync(job_id: str) -> bool:
"""Запросить отмену и прервать активный подпроцесс (если есть)."""
with _thread_lock:
ev = _cancel_events.get(job_id)
if ev is None:
return False
ev.set()
logger.info("Запрошена отмена задания %s", job_id)
kill_subprocess_for_job_sync(job_id)
return True
def is_cancelled_sync(job_id: str) -> bool:
"""Проверка из worker-thread между этапами и перед каждой следующей командой."""
with _thread_lock:
ev = _cancel_events.get(job_id)
return bool(ev and ev.is_set())
@dataclass
class JobRecord:
"""Описание одного задания."""
job_id: str
kind: str
status: JobStatus
cluster_name: str | None
created_at_utc: str
message: str | None = None
result: dict[str, Any] | None = None
# Журнал после завершения (stdout/stderr kind create и этапы); пока задание активно — см. deque
log_lines: list[str] = field(default_factory=list)
def _record_to_dict(rec: JobRecord) -> dict[str, Any]:
return {
"job_id": rec.job_id,
"kind": rec.kind,
"status": rec.status,
"cluster_name": rec.cluster_name,
"created_at_utc": rec.created_at_utc,
"message": rec.message,
"result": rec.result,
"log_lines": list(rec.log_lines),
}
def _dict_to_record(d: Any) -> JobRecord | None:
if not isinstance(d, dict):
return None
jid = d.get("job_id")
kind = d.get("kind")
status = d.get("status")
if not isinstance(jid, str) or not jid or not isinstance(kind, str) or not kind.strip():
return None
if status not in ("queued", "running", "success", "failed", "cancelled"):
return None
created = d.get("created_at_utc")
if not isinstance(created, str) or not created:
return None
cn = d.get("cluster_name")
if cn is None:
cluster_name = None
elif isinstance(cn, str):
cluster_name = cn
else:
cluster_name = str(cn)
msg = d.get("message")
if msg is None:
message = None
elif isinstance(msg, str):
message = msg
else:
message = str(msg)
result = d.get("result")
if result is not None and not isinstance(result, dict):
result = None
lines = d.get("log_lines")
log_lines: list[str] = []
if isinstance(lines, list):
log_lines = [str(x) for x in lines]
return JobRecord(
job_id=jid,
kind=kind,
status=cast(JobStatus, status),
cluster_name=cluster_name,
created_at_utc=created,
message=message,
result=result,
log_lines=log_lines,
)
def _atomic_write_jobs_file(path: Path, jobs: list[dict[str, Any]]) -> None:
"""Запись JSON атомарно (temp + replace)."""
path.parent.mkdir(parents=True, exist_ok=True)
payload = {"version": _JOBS_FILE_VERSION, "jobs": jobs}
tmp = path.with_suffix(path.suffix + ".tmp")
text = json.dumps(payload, ensure_ascii=False, indent=2, default=str)
tmp.write_text(text, encoding="utf-8")
tmp.replace(path)
class JobStore:
"""Потокобезопасное (asyncio) хранилище заданий с персистентностью в JSON."""
def __init__(self) -> None:
self._jobs: dict[str, JobRecord] = {}
self._lock = asyncio.Lock()
self._load_from_disk_sync()
def _load_from_disk_sync(self) -> None:
path = _jobs_file_path()
if not path.is_file():
logger.debug("Файл истории заданий отсутствует: %s", path)
return
try:
raw = path.read_text(encoding="utf-8")
data = json.loads(raw)
except (OSError, json.JSONDecodeError) as e:
logger.warning("Не удалось прочитать историю заданий %s: %s", path, e)
return
jobs_raw = data.get("jobs") if isinstance(data, dict) else data
if not isinstance(jobs_raw, list):
logger.warning("Некорректный формат истории заданий (ожидался ключ jobs или массив)")
return
interrupt_msg = "Прервано перезапуском сервиса (задание не было завершено)."
loaded: list[JobRecord] = []
for item in jobs_raw:
rec = _dict_to_record(item)
if rec is None:
continue
if rec.status in ("queued", "running"):
rec = JobRecord(
job_id=rec.job_id,
kind=rec.kind,
status="failed",
cluster_name=rec.cluster_name,
created_at_utc=rec.created_at_utc,
message=interrupt_msg,
result=rec.result,
log_lines=rec.log_lines,
)
loaded.append(rec)
loaded.sort(key=lambda r: r.created_at_utc)
if len(loaded) > _MAX_JOBS:
loaded = loaded[-_MAX_JOBS:]
self._jobs = {r.job_id: r for r in loaded}
logger.info("Восстановлено заданий из %s: %s", path, len(self._jobs))
async def _persist_async(self) -> None:
"""Сохранить текущий снимок заданий на диск (вне удержания lock во время записи)."""
async with self._lock:
items = sorted(self._jobs.values(), key=lambda r: r.created_at_utc)
snapshot = [_record_to_dict(r) for r in items]
path = _jobs_file_path()
try:
await asyncio.to_thread(_atomic_write_jobs_file, path, snapshot)
logger.debug("История заданий записана: %s (%s записей)", path, len(snapshot))
except OSError as e:
logger.error("Не удалось сохранить историю заданий в %s: %s", path, e)
async def create_job(self, kind: str, *, cluster_name: str | None) -> JobRecord:
"""Зарегистрировать задание в статусе ``queued``."""
jid = uuid.uuid4().hex
now = datetime.now(timezone.utc).isoformat()
rec = JobRecord(
job_id=jid,
kind=kind,
status="queued",
cluster_name=cluster_name,
created_at_utc=now,
)
async with self._lock:
self._jobs[jid] = rec
while len(self._jobs) > _MAX_JOBS:
oldest_id = min(self._jobs, key=lambda k: self._jobs[k].created_at_utc)
end_job_tracking(oldest_id)
del self._jobs[oldest_id]
logger.debug("Вытеснено старое задание из хранилища: %s", oldest_id)
begin_job_tracking(jid)
logger.info("Создано задание %s kind=%s cluster=%s", jid, kind, cluster_name)
await self._persist_async()
return rec
async def set_running(
self,
job_id: str,
*,
stage: str = "Выполняется…",
percent: int = 5,
) -> None:
async with self._lock:
if job_id in self._jobs:
self._jobs[job_id].status = "running"
self._jobs[job_id].message = None
set_progress_sync(job_id, stage, percent)
await self._persist_async()
async def set_success(self, job_id: str, *, result: dict[str, Any] | None = None, message: str | None = None) -> None:
logs = take_logs_finalize_sync(job_id)
async with self._lock:
if job_id in self._jobs:
self._jobs[job_id].status = "success"
self._jobs[job_id].result = result
self._jobs[job_id].message = message
self._jobs[job_id].log_lines = logs
set_progress_sync(job_id, "Готово", 100)
await self._persist_async()
async def set_failed(self, job_id: str, message: str) -> None:
logs = take_logs_finalize_sync(job_id)
async with self._lock:
if job_id in self._jobs:
self._jobs[job_id].status = "failed"
self._jobs[job_id].message = message
self._jobs[job_id].log_lines = logs
logger.warning("Задание %s завершилось ошибкой: %s", job_id, message)
await self._persist_async()
async def set_cancelled(self, job_id: str, message: str = "Операция отменена") -> None:
logs = take_logs_finalize_sync(job_id)
async with self._lock:
if job_id in self._jobs:
self._jobs[job_id].status = "cancelled"
self._jobs[job_id].message = message
self._jobs[job_id].log_lines = logs
logger.info("Задание %s отменено: %s", job_id, message)
await self._persist_async()
async def get(self, job_id: str) -> JobRecord | None:
async with self._lock:
return self._jobs.get(job_id)
def snapshot_all(self) -> list[JobRecord]:
"""Снимок всех заданий (для отладки; без блокировки — eventual consistency)."""
return list(self._jobs.values())
def snapshot_recent_sorted(self, *, limit: int) -> list[JobRecord]:
"""Задания от новых к старым, не более ``limit``."""
items = self.snapshot_all()
items.sort(key=lambda r: r.created_at_utc, reverse=True)
return items[: max(1, limit)]
async def purge_completed(self) -> int:
"""
Удалить из памяти только завершённые задания (success / failed / cancelled).
Активные (queued, running) не трогаем. Служебные структуры отслеживания подчищаются.
"""
async with self._lock:
to_remove = [
jid
for jid, rec in self._jobs.items()
if rec.status in ("success", "failed", "cancelled")
]
for jid in to_remove:
del self._jobs[jid]
for jid in to_remove:
end_job_tracking(jid)
n = len(to_remove)
await self._persist_async()
return n
# Синглтон на процесс uvicorn
job_store = JobStore()