6f3daa33ec
- Шапка: логотип Kubernetes, ссылка на главную, выпадающее меню API (Swagger/ReDoc/Health), переключатель светлой/тёмной темы (localStorage). - Светлая тема в синей гамме; выравнивание кнопки темы в ряду с пилюлями. - Дашборд: единая карточка ошибки health/stats, подсказка Docker/Podman, поле container_cli в GET /stats, total_workers_from_meta всегда число (0 без meta). - Правки кластеров, job_store, compose, документация и частичные шаблоны.
436 lines
17 KiB
Python
436 lines
17 KiB
Python
"""Хранилище фоновых заданий (создание/старт/стоп кластера).
|
||
|
||
Снимок заданий сохраняется в 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]] = {}
|
||
# Активный дочерний процесс (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 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:
|
||
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)
|
||
|
||
|
||
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()
|