Журнал: пагинация 30 записей на страницу, API offset/total_pages
- job_journal: collect_recent_journal_entries_page_sync(limit, offset) - GET /api/v1/journal/recent: limit по умолчанию 30, offset, total, page, total_pages - journal.html/js: навигация Первая/Назад/номера/Вперёд/Последняя, стили - app/docs/api_routes.md: описание query и пример ответа - Прочие изменения UI/API (аддоны, helm, job_journal в кластерах) в том же коммите
This commit is contained in:
@@ -0,0 +1,512 @@
|
||||
"""Установка и удаление типовых Helm-чартов в выбранном кластере (kubeconfig с диска).
|
||||
|
||||
Используются бинарники ``helm`` и ``kubectl`` из PATH; ``KUBECONFIG`` — путь из
|
||||
:func:`kubeconfig_patch.kubeconfig_path_for_container_kubectl` (как для kubectl в приложении).
|
||||
|
||||
Автор: Сергей Антропов
|
||||
Сайт: https://devops.org.ru
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import shutil
|
||||
import subprocess
|
||||
import tempfile
|
||||
import threading
|
||||
import time
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
import yaml
|
||||
|
||||
from core.cluster_lifecycle import validate_cluster_name
|
||||
from kind_k8s_paths import clusters_dir
|
||||
from kubeconfig_patch import kubeconfig_path_for_container_kubectl
|
||||
|
||||
logger = logging.getLogger("kind_k8s.helm_addons")
|
||||
|
||||
# Официальные репозитории чартов (имя → URL).
|
||||
HELM_REPOS: dict[str, str] = {
|
||||
"ingress-nginx": "https://kubernetes.github.io/ingress-nginx",
|
||||
"prometheus-community": "https://prometheus-community.github.io/helm-charts",
|
||||
"metrics-server": "https://kubernetes-sigs.github.io/metrics-server/",
|
||||
"istio": "https://istio-release.storage.googleapis.com/charts",
|
||||
"kiali": "https://kiali.org/helm-charts",
|
||||
}
|
||||
|
||||
# Ссылки для ``helm search repo <ref> --versions`` (ключи совпадают с полями API версий чартов).
|
||||
HELM_ADDON_CHART_REFS: dict[str, str] = {
|
||||
"ingress_nginx": "ingress-nginx/ingress-nginx",
|
||||
"kube_prometheus_stack": "prometheus-community/kube-prometheus-stack",
|
||||
"metrics_server": "metrics-server/metrics-server",
|
||||
# Для base и istiod обычно берут одну версию чарта Istio.
|
||||
"istio": "istio/base",
|
||||
"kiali_server": "kiali/kiali-server",
|
||||
}
|
||||
|
||||
_versions_cache: dict[str, tuple[float, list[str]]] = {}
|
||||
_versions_lock = threading.Lock()
|
||||
|
||||
|
||||
class HelmAddonError(Exception):
|
||||
"""Ошибка helm/kubectl (сообщение для API)."""
|
||||
|
||||
def __init__(self, message: str, *, exit_code: int = 1) -> None:
|
||||
super().__init__(message)
|
||||
self.exit_code = exit_code
|
||||
|
||||
|
||||
def _helm_timeout_sec() -> int:
|
||||
raw = (os.environ.get("KIND_K8S_HELM_TIMEOUT_SEC") or "900").strip()
|
||||
try:
|
||||
return max(60, min(int(raw), 7200))
|
||||
except ValueError:
|
||||
return 900
|
||||
|
||||
|
||||
def _helm_versions_cache_sec() -> int:
|
||||
"""TTL кэша списков версий (``helm search repo`` после ``repo update``)."""
|
||||
raw = (os.environ.get("KIND_K8S_HELM_VERSIONS_CACHE_SEC") or "600").strip()
|
||||
try:
|
||||
return max(30, min(int(raw), 86400))
|
||||
except ValueError:
|
||||
return 600
|
||||
|
||||
|
||||
def _helm_versions_max_list() -> int:
|
||||
raw = (os.environ.get("KIND_K8S_HELM_VERSIONS_MAX") or "80").strip()
|
||||
try:
|
||||
return max(10, min(int(raw), 200))
|
||||
except ValueError:
|
||||
return 80
|
||||
|
||||
|
||||
def search_repo_chart_versions(chart_ref: str) -> list[str]:
|
||||
"""
|
||||
Уникальные версии чарта (новые первыми), через ``helm search repo … --versions -o json``.
|
||||
|
||||
Перед вызовом рекомендуется :func:`helm_repo_ensure`; здесь кэш по ``chart_ref``.
|
||||
"""
|
||||
ref = chart_ref.strip()
|
||||
if not ref:
|
||||
return []
|
||||
now = time.monotonic()
|
||||
ttl = _helm_versions_cache_sec()
|
||||
with _versions_lock:
|
||||
hit = _versions_cache.get(ref)
|
||||
if hit is not None:
|
||||
ts, vers = hit
|
||||
if now - ts < ttl:
|
||||
logger.debug("Версии чарта %s из кэша (%s шт.)", ref, len(vers))
|
||||
return list(vers)
|
||||
|
||||
require_helm_binary()
|
||||
rc, out, err = _run_cmd(
|
||||
["helm", "search", "repo", ref, "--versions", "-o", "json"],
|
||||
timeout=180,
|
||||
)
|
||||
if rc != 0:
|
||||
logger.warning("helm search repo %s: %s", ref, err or out or rc)
|
||||
with _versions_lock:
|
||||
_versions_cache[ref] = (now, [])
|
||||
return []
|
||||
|
||||
versions: list[str] = []
|
||||
seen: set[str] = set()
|
||||
try:
|
||||
rows = json.loads(out or "[]")
|
||||
except json.JSONDecodeError:
|
||||
rows = []
|
||||
cap = _helm_versions_max_list()
|
||||
if isinstance(rows, list):
|
||||
for row in rows:
|
||||
if not isinstance(row, dict):
|
||||
continue
|
||||
v = str(row.get("version", "")).strip()
|
||||
if not v or v in seen:
|
||||
continue
|
||||
seen.add(v)
|
||||
versions.append(v)
|
||||
if len(versions) >= cap:
|
||||
break
|
||||
|
||||
with _versions_lock:
|
||||
_versions_cache[ref] = (now, versions)
|
||||
logger.info("Получены версии для %s: %s записей", ref, len(versions))
|
||||
return versions
|
||||
|
||||
|
||||
def helm_addon_chart_versions_all() -> dict[str, list[str]]:
|
||||
"""
|
||||
Версии для всех аддонов UI (после обновления репозиториев — один раз).
|
||||
|
||||
Ключи: ``ingress_nginx``, ``kube_prometheus_stack``, ``metrics_server``, ``istio``, ``kiali_server``.
|
||||
"""
|
||||
helm_repo_ensure()
|
||||
out: dict[str, list[str]] = {}
|
||||
for key, cref in HELM_ADDON_CHART_REFS.items():
|
||||
out[key] = search_repo_chart_versions(cref)
|
||||
return out
|
||||
|
||||
|
||||
def require_helm_binary() -> None:
|
||||
if not shutil.which("helm"):
|
||||
raise HelmAddonError("В образе не найден helm в PATH. Пересоберите образ с установленным Helm.", exit_code=127)
|
||||
|
||||
|
||||
def require_kubectl_binary() -> None:
|
||||
if not shutil.which("kubectl"):
|
||||
raise HelmAddonError("Не найден kubectl в PATH.", exit_code=127)
|
||||
|
||||
|
||||
def _kubeconfig_runtime_path(cluster_name: str) -> Path:
|
||||
if not validate_cluster_name(cluster_name):
|
||||
raise HelmAddonError("Некорректное имя кластера")
|
||||
src = clusters_dir() / cluster_name / "kubeconfig"
|
||||
if not src.is_file():
|
||||
raise HelmAddonError(f"Нет kubeconfig: {src}")
|
||||
return kubeconfig_path_for_container_kubectl(cluster_name=cluster_name, kube_source_path=src)
|
||||
|
||||
|
||||
def _run_cmd(
|
||||
cmd: list[str],
|
||||
*,
|
||||
timeout: int | None = None,
|
||||
env: dict[str, str] | None = None,
|
||||
) -> tuple[int, str, str]:
|
||||
"""Выполнить команду; возвращает (код, stdout, stderr)."""
|
||||
t = timeout if timeout is not None else _helm_timeout_sec()
|
||||
logger.info("Выполнение: %s (timeout=%s)", " ".join(cmd[:12]) + (" …" if len(cmd) > 12 else ""), t)
|
||||
try:
|
||||
p = subprocess.run(
|
||||
cmd,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=t,
|
||||
env={**os.environ, **(env or {})},
|
||||
)
|
||||
except subprocess.TimeoutExpired:
|
||||
logger.warning("Таймаут команды: %s", cmd[0])
|
||||
raise HelmAddonError(f"Превышено время ожидания ({t} с) для команды «{cmd[0]}»", exit_code=124) from None
|
||||
out = (p.stdout or "").strip()
|
||||
err = (p.stderr or "").strip()
|
||||
return p.returncode, out, err
|
||||
|
||||
|
||||
def helm_repo_ensure() -> tuple[int, str]:
|
||||
"""Добавить недостающие репозитории и выполнить ``helm repo update``."""
|
||||
require_helm_binary()
|
||||
lines: list[str] = []
|
||||
for name, url in HELM_REPOS.items():
|
||||
rc, o, e = _run_cmd(["helm", "repo", "add", name, url], timeout=120)
|
||||
if rc != 0 and "already exists" not in (e + o).lower():
|
||||
raise HelmAddonError(f"helm repo add {name}: {e or o or rc}", exit_code=rc)
|
||||
if o:
|
||||
lines.append(o)
|
||||
if e and "already exists" not in e.lower():
|
||||
lines.append(e)
|
||||
rc2, o2, e2 = _run_cmd(["helm", "repo", "update"], timeout=300)
|
||||
if rc2 != 0:
|
||||
raise HelmAddonError(f"helm repo update: {e2 or o2 or rc2}", exit_code=rc2)
|
||||
lines.append(o2 or e2)
|
||||
return 0, "\n".join(x for x in lines if x).strip()
|
||||
|
||||
|
||||
def helm_list_all(cluster_name: str) -> list[dict[str, Any]]:
|
||||
"""Список релизов ``helm list -A -o json``."""
|
||||
kcfg = str(_kubeconfig_runtime_path(cluster_name))
|
||||
rc, out, err = _run_cmd(
|
||||
["helm", "list", "-A", "-o", "json", "--kubeconfig", kcfg],
|
||||
timeout=120,
|
||||
)
|
||||
if rc != 0:
|
||||
logger.warning("helm list: %s", err or out)
|
||||
return []
|
||||
if not out.strip():
|
||||
return []
|
||||
try:
|
||||
data = json.loads(out)
|
||||
return data if isinstance(data, list) else []
|
||||
except json.JSONDecodeError:
|
||||
return []
|
||||
|
||||
|
||||
def addon_status(cluster_name: str) -> dict[str, bool]:
|
||||
"""Какие аддоны установлены (по имени релиза и namespace)."""
|
||||
lst = helm_list_all(cluster_name)
|
||||
names_ns = {(str(x.get("name", "")), str(x.get("namespace", ""))) for x in lst if isinstance(x, dict)}
|
||||
return {
|
||||
"ingress_nginx": ("ingress-nginx", "ingress-nginx") in names_ns,
|
||||
"kube_prometheus_stack": ("kube-prometheus-stack", "monitoring") in names_ns,
|
||||
"metrics_server": ("metrics-server", "kube-system") in names_ns,
|
||||
"istio_base": ("istio-base", "istio-system") in names_ns,
|
||||
"istiod": ("istiod", "istio-system") in names_ns,
|
||||
"kiali": ("kiali-server", "istio-system") in names_ns,
|
||||
}
|
||||
|
||||
|
||||
def _kubectl_apply_manifest(cluster_name: str, manifest: dict[str, Any]) -> tuple[int, str]:
|
||||
require_kubectl_binary()
|
||||
kcfg = str(_kubeconfig_runtime_path(cluster_name))
|
||||
yml = yaml.safe_dump(manifest, allow_unicode=True, default_flow_style=False)
|
||||
p = subprocess.run(
|
||||
["kubectl", "apply", "-f", "-", "--kubeconfig", kcfg, "--request-timeout=120s"],
|
||||
input=yml,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=180,
|
||||
)
|
||||
msg = ((p.stdout or "").strip() + "\n" + (p.stderr or "").strip()).strip()
|
||||
return p.returncode, msg
|
||||
|
||||
|
||||
def _ensure_namespace(cluster_name: str, namespace: str) -> None:
|
||||
manifest = {
|
||||
"apiVersion": "v1",
|
||||
"kind": "Namespace",
|
||||
"metadata": {"name": namespace},
|
||||
}
|
||||
rc, msg = _kubectl_apply_manifest(cluster_name, manifest)
|
||||
if rc != 0:
|
||||
raise HelmAddonError(f"Создание namespace {namespace}: {msg}", exit_code=rc)
|
||||
|
||||
|
||||
def _helm_upgrade_install(
|
||||
cluster_name: str,
|
||||
release: str,
|
||||
chart: str,
|
||||
namespace: str,
|
||||
*,
|
||||
chart_version: str | None = None,
|
||||
extra_args: list[str] | None = None,
|
||||
values_file: Path | None = None,
|
||||
) -> tuple[int, str, str]:
|
||||
kcfg = str(_kubeconfig_runtime_path(cluster_name))
|
||||
cmd = [
|
||||
"helm",
|
||||
"upgrade",
|
||||
"--install",
|
||||
release,
|
||||
chart,
|
||||
"-n",
|
||||
namespace,
|
||||
"--create-namespace",
|
||||
"--kubeconfig",
|
||||
kcfg,
|
||||
"--wait",
|
||||
"--timeout",
|
||||
f"{_helm_timeout_sec()}s",
|
||||
]
|
||||
ver = (chart_version or "").strip()
|
||||
if ver:
|
||||
cmd.extend(["--version", ver])
|
||||
if values_file is not None:
|
||||
cmd.extend(["-f", str(values_file)])
|
||||
if extra_args:
|
||||
cmd.extend(extra_args)
|
||||
return _run_cmd(cmd, timeout=_helm_timeout_sec() + 60)
|
||||
|
||||
|
||||
def _helm_uninstall(cluster_name: str, release: str, namespace: str) -> tuple[int, str, str]:
|
||||
kcfg = str(_kubeconfig_runtime_path(cluster_name))
|
||||
return _run_cmd(
|
||||
["helm", "uninstall", release, "-n", namespace, "--kubeconfig", kcfg, "--wait", "--timeout", "600s"],
|
||||
timeout=660,
|
||||
)
|
||||
|
||||
|
||||
# --- Установка / удаление по типам ---
|
||||
|
||||
|
||||
def install_ingress_nginx(cluster_name: str, chart_version: str | None) -> tuple[bool, str]:
|
||||
"""ingress-nginx с NodePort для kind."""
|
||||
helm_repo_ensure()
|
||||
extra: list[str] = [
|
||||
"--set",
|
||||
"controller.service.type=NodePort",
|
||||
"--set",
|
||||
"controller.service.nodePorts.http=30080",
|
||||
]
|
||||
rc, out, err = _helm_upgrade_install(
|
||||
cluster_name,
|
||||
"ingress-nginx",
|
||||
"ingress-nginx/ingress-nginx",
|
||||
"ingress-nginx",
|
||||
chart_version=chart_version,
|
||||
extra_args=extra,
|
||||
)
|
||||
text = "\n".join(filter(None, [out, err])).strip()
|
||||
if rc != 0:
|
||||
raise HelmAddonError(text or f"helm exit {rc}", exit_code=rc)
|
||||
return True, text or "ingress-nginx установлен."
|
||||
|
||||
|
||||
def uninstall_ingress_nginx(cluster_name: str) -> tuple[bool, str]:
|
||||
rc, out, err = _helm_uninstall(cluster_name, "ingress-nginx", "ingress-nginx")
|
||||
text = "\n".join(filter(None, [out, err])).strip()
|
||||
if rc != 0 and "not found" not in text.lower():
|
||||
raise HelmAddonError(text or f"helm exit {rc}", exit_code=rc)
|
||||
return True, text or "ingress-nginx удалён (или отсутствовал)."
|
||||
|
||||
|
||||
def install_kube_prometheus_stack(
|
||||
cluster_name: str,
|
||||
*,
|
||||
grafana_admin_user: str,
|
||||
grafana_admin_password: str,
|
||||
chart_version: str | None = None,
|
||||
) -> tuple[bool, str]:
|
||||
helm_repo_ensure()
|
||||
values = {
|
||||
"grafana": {
|
||||
"adminUser": grafana_admin_user,
|
||||
"adminPassword": grafana_admin_password,
|
||||
}
|
||||
}
|
||||
with tempfile.NamedTemporaryFile(mode="w", suffix=".yaml", delete=False, encoding="utf-8") as tmp:
|
||||
yaml.safe_dump(values, tmp, allow_unicode=True)
|
||||
tmp_path = Path(tmp.name)
|
||||
try:
|
||||
rc, out, err = _helm_upgrade_install(
|
||||
cluster_name,
|
||||
"kube-prometheus-stack",
|
||||
"prometheus-community/kube-prometheus-stack",
|
||||
"monitoring",
|
||||
chart_version=chart_version,
|
||||
values_file=tmp_path,
|
||||
)
|
||||
finally:
|
||||
tmp_path.unlink(missing_ok=True)
|
||||
text = "\n".join(filter(None, [out, err])).strip()
|
||||
if rc != 0:
|
||||
raise HelmAddonError(text or f"helm exit {rc}", exit_code=rc)
|
||||
return True, text or "kube-prometheus-stack установлен."
|
||||
|
||||
|
||||
def uninstall_kube_prometheus_stack(cluster_name: str) -> tuple[bool, str]:
|
||||
rc, out, err = _helm_uninstall(cluster_name, "kube-prometheus-stack", "monitoring")
|
||||
text = "\n".join(filter(None, [out, err])).strip()
|
||||
if rc != 0 and "not found" not in text.lower():
|
||||
raise HelmAddonError(text or f"helm exit {rc}", exit_code=rc)
|
||||
return True, text or "kube-prometheus-stack удалён (или отсутствовал)."
|
||||
|
||||
|
||||
def install_metrics_server(cluster_name: str, chart_version: str | None = None) -> tuple[bool, str]:
|
||||
helm_repo_ensure()
|
||||
values = {
|
||||
"args": [
|
||||
"--kubelet-insecure-tls",
|
||||
"--kubelet-preferred-address-types=InternalIP,Hostname,ExternalIP",
|
||||
]
|
||||
}
|
||||
with tempfile.NamedTemporaryFile(mode="w", suffix=".yaml", delete=False, encoding="utf-8") as tmp:
|
||||
yaml.safe_dump(values, tmp, allow_unicode=True)
|
||||
tmp_path = Path(tmp.name)
|
||||
try:
|
||||
rc, out, err = _helm_upgrade_install(
|
||||
cluster_name,
|
||||
"metrics-server",
|
||||
"metrics-server/metrics-server",
|
||||
"kube-system",
|
||||
chart_version=chart_version,
|
||||
values_file=tmp_path,
|
||||
)
|
||||
finally:
|
||||
tmp_path.unlink(missing_ok=True)
|
||||
text = "\n".join(filter(None, [out, err])).strip()
|
||||
if rc != 0:
|
||||
raise HelmAddonError(text or f"helm exit {rc}", exit_code=rc)
|
||||
return True, text or "metrics-server установлен."
|
||||
|
||||
|
||||
def uninstall_metrics_server(cluster_name: str) -> tuple[bool, str]:
|
||||
rc, out, err = _helm_uninstall(cluster_name, "metrics-server", "kube-system")
|
||||
text = "\n".join(filter(None, [out, err])).strip()
|
||||
if rc != 0 and "not found" not in text.lower():
|
||||
raise HelmAddonError(text or f"helm exit {rc}", exit_code=rc)
|
||||
return True, text or "metrics-server удалён (или отсутствовал)."
|
||||
|
||||
|
||||
def install_istio_and_kiali(
|
||||
cluster_name: str,
|
||||
*,
|
||||
kiali_username: str,
|
||||
kiali_password: str,
|
||||
istio_chart_version: str | None = None,
|
||||
kiali_chart_version: str | None = None,
|
||||
) -> tuple[bool, str]:
|
||||
"""Istio base + istiod + секрет входа Kiali + chart kiali-server (strategy=login)."""
|
||||
helm_repo_ensure()
|
||||
_ensure_namespace(cluster_name, "istio-system")
|
||||
|
||||
log_parts: list[str] = []
|
||||
istio_ver = (istio_chart_version or "").strip() or None
|
||||
kiali_ver = (kiali_chart_version or "").strip() or None
|
||||
|
||||
rc, out, err = _helm_upgrade_install(
|
||||
cluster_name,
|
||||
"istio-base",
|
||||
"istio/base",
|
||||
"istio-system",
|
||||
chart_version=istio_ver,
|
||||
)
|
||||
log_parts.append(out or err)
|
||||
if rc != 0:
|
||||
raise HelmAddonError((err or out or f"istio-base exit {rc}")[:8000], exit_code=rc)
|
||||
|
||||
rc2, o2, e2 = _helm_upgrade_install(
|
||||
cluster_name,
|
||||
"istiod",
|
||||
"istio/istiod",
|
||||
"istio-system",
|
||||
chart_version=istio_ver,
|
||||
)
|
||||
log_parts.append(o2 or e2)
|
||||
if rc2 != 0:
|
||||
raise HelmAddonError((e2 or o2 or f"istiod exit {rc2}")[:8000], exit_code=rc2)
|
||||
|
||||
# Секрет для стратегии login (ключи username / passphrase — ожидание Kiali).
|
||||
secret = {
|
||||
"apiVersion": "v1",
|
||||
"kind": "Secret",
|
||||
"metadata": {"name": "kiali", "namespace": "istio-system"},
|
||||
"type": "Opaque",
|
||||
"stringData": {
|
||||
"username": kiali_username,
|
||||
"passphrase": kiali_password,
|
||||
},
|
||||
}
|
||||
rcs, msgs = _kubectl_apply_manifest(cluster_name, secret)
|
||||
if rcs != 0:
|
||||
raise HelmAddonError(f"Secret kiali: {msgs}", exit_code=rcs)
|
||||
log_parts.append(msgs)
|
||||
|
||||
rc3, o3, e3 = _helm_upgrade_install(
|
||||
cluster_name,
|
||||
"kiali-server",
|
||||
"kiali/kiali-server",
|
||||
"istio-system",
|
||||
chart_version=kiali_ver,
|
||||
extra_args=["--set", "auth.strategy=login"],
|
||||
)
|
||||
log_parts.append(o3 or e3)
|
||||
if rc3 != 0:
|
||||
raise HelmAddonError((e3 or o3 or f"kiali-server exit {rc3}")[:8000], exit_code=rc3)
|
||||
|
||||
return True, "\n\n".join(x for x in log_parts if x).strip() or "Istio и Kiali установлены."
|
||||
|
||||
|
||||
def uninstall_istio_and_kiali(cluster_name: str) -> tuple[bool, str]:
|
||||
"""Обратный порядок: Kiali → istiod → istio-base."""
|
||||
parts: list[str] = []
|
||||
for rel in ("kiali-server", "istiod", "istio-base"):
|
||||
rc, out, err = _helm_uninstall(cluster_name, rel, "istio-system")
|
||||
line = "\n".join(filter(None, [out, err])).strip()
|
||||
if line:
|
||||
parts.append(line)
|
||||
if rc != 0 and "not found" not in line.lower() and "release: not found" not in line.lower():
|
||||
raise HelmAddonError(line or f"helm uninstall {rel} код {rc}", exit_code=rc)
|
||||
return True, "\n".join(parts).strip() or "Istio/Kiali удалены (или отсутствовали)."
|
||||
@@ -0,0 +1,205 @@
|
||||
"""История завершённых заданий (create/start/stop и т.д.) в ``clusters/<имя>/journal/jobs_history.json``.
|
||||
|
||||
Дополняет глобальный ``kind_k8s_jobs.json``: записи дублируются в каталог кластера для архива на томе.
|
||||
|
||||
Автор: Сергей Антропов
|
||||
Сайт: https://devops.org.ru
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import threading
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from core.cluster_lifecycle import validate_cluster_name
|
||||
from kind_k8s_paths import clusters_dir
|
||||
|
||||
logger = logging.getLogger("kind_k8s.job_journal")
|
||||
|
||||
JOURNAL_SUBDIR = "journal"
|
||||
JOBS_HISTORY_FILENAME = "jobs_history.json"
|
||||
JOURNAL_FILE_VERSION = 1
|
||||
|
||||
# Блокировки по имени кластера (один процесс uvicorn).
|
||||
_locks_guard = threading.Lock()
|
||||
_cluster_locks: dict[str, threading.Lock] = {}
|
||||
|
||||
|
||||
def _cluster_lock(name: str) -> threading.Lock:
|
||||
with _locks_guard:
|
||||
if name not in _cluster_locks:
|
||||
_cluster_locks[name] = threading.Lock()
|
||||
return _cluster_locks[name]
|
||||
|
||||
|
||||
def _max_entries() -> int:
|
||||
raw = (os.environ.get("KIND_K8S_CLUSTER_JOURNAL_MAX_ENTRIES") or "500").strip()
|
||||
try:
|
||||
return max(20, min(int(raw), 5000))
|
||||
except ValueError:
|
||||
return 500
|
||||
|
||||
|
||||
def _max_log_lines() -> int:
|
||||
raw = (os.environ.get("KIND_K8S_CLUSTER_JOURNAL_MAX_LOG_LINES") or "2000").strip()
|
||||
try:
|
||||
return max(50, min(int(raw), 20000))
|
||||
except ValueError:
|
||||
return 2000
|
||||
|
||||
|
||||
def journal_file_path(cluster_name: str) -> Path:
|
||||
"""Путь к JSON с историей заданий кластера."""
|
||||
return clusters_dir() / cluster_name.strip() / JOURNAL_SUBDIR / JOBS_HISTORY_FILENAME
|
||||
|
||||
|
||||
def append_job_from_record_sync(
|
||||
*,
|
||||
job_id: str,
|
||||
kind: str,
|
||||
cluster_name: str | None,
|
||||
status: str,
|
||||
message: str | None,
|
||||
created_at_utc: str,
|
||||
log_lines: list[str],
|
||||
result: dict[str, Any] | None,
|
||||
) -> None:
|
||||
"""
|
||||
Добавить запись о завершённом задании в журнал каталога кластера.
|
||||
|
||||
Вызывается из потока после ``job_store`` (to_thread); без паролей в result.
|
||||
"""
|
||||
if not cluster_name or not str(cluster_name).strip():
|
||||
return
|
||||
name = str(cluster_name).strip()
|
||||
if not validate_cluster_name(name):
|
||||
logger.debug("Журнал кластера: пропуск, некорректное имя «%s»", name)
|
||||
return
|
||||
|
||||
cdir = clusters_dir() / name
|
||||
if not cdir.is_dir():
|
||||
logger.debug("Журнал кластера: каталог %s отсутствует", cdir)
|
||||
return
|
||||
|
||||
jdir = cdir / JOURNAL_SUBDIR
|
||||
finished = datetime.now(timezone.utc).isoformat()
|
||||
cap_lines = _max_log_lines()
|
||||
lines = list(log_lines)[-cap_lines:] if log_lines else []
|
||||
|
||||
entry: dict[str, Any] = {
|
||||
"version": JOURNAL_FILE_VERSION,
|
||||
"job_id": job_id,
|
||||
"kind": kind,
|
||||
"cluster_name": name,
|
||||
"status": status,
|
||||
"message": message,
|
||||
"created_at_utc": created_at_utc,
|
||||
"finished_at_utc": finished,
|
||||
"log_lines": lines,
|
||||
"result": result,
|
||||
}
|
||||
|
||||
path = jdir / JOBS_HISTORY_FILENAME
|
||||
lock = _cluster_lock(name)
|
||||
with lock:
|
||||
jdir.mkdir(parents=True, exist_ok=True)
|
||||
entries: list[dict[str, Any]] = []
|
||||
if path.is_file():
|
||||
try:
|
||||
raw = path.read_text(encoding="utf-8")
|
||||
data = json.loads(raw)
|
||||
if isinstance(data, dict) and isinstance(data.get("entries"), list):
|
||||
entries = [e for e in data["entries"] if isinstance(e, dict)]
|
||||
except (OSError, json.JSONDecodeError) as e:
|
||||
logger.warning("Журнал %s: не прочитан, создаём заново: %s", path, e)
|
||||
entries = []
|
||||
|
||||
entries.insert(0, entry)
|
||||
max_e = _max_entries()
|
||||
if len(entries) > max_e:
|
||||
entries = entries[:max_e]
|
||||
|
||||
payload = {"file_version": JOURNAL_FILE_VERSION, "entries": entries}
|
||||
tmp = path.with_suffix(path.suffix + ".tmp")
|
||||
try:
|
||||
tmp.write_text(json.dumps(payload, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
|
||||
tmp.replace(path)
|
||||
logger.info("Журнал кластера «%s»: запись %s (%s)", name, job_id, kind)
|
||||
except OSError as e:
|
||||
logger.warning("Журнал %s: запись не удалась: %s", path, e)
|
||||
try:
|
||||
tmp.unlink(missing_ok=True)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
def read_cluster_journal_sync(cluster_name: str) -> dict[str, Any] | None:
|
||||
"""Прочитать файл журнала кластера или ``None``."""
|
||||
n = cluster_name.strip()
|
||||
if not validate_cluster_name(n):
|
||||
return None
|
||||
path = journal_file_path(n)
|
||||
if not path.is_file():
|
||||
return None
|
||||
try:
|
||||
return json.loads(path.read_text(encoding="utf-8"))
|
||||
except (OSError, json.JSONDecodeError):
|
||||
return None
|
||||
|
||||
|
||||
def _collect_all_journal_rows_sorted_sync() -> list[tuple[str, dict[str, Any]]]:
|
||||
"""Все записи из журналов кластеров, отсортированные по времени (новые первыми)."""
|
||||
root = clusters_dir()
|
||||
if not root.is_dir():
|
||||
return []
|
||||
flat: list[tuple[str, dict[str, Any]]] = []
|
||||
for sub in sorted(root.iterdir()):
|
||||
if not sub.is_dir():
|
||||
continue
|
||||
name = sub.name
|
||||
if not validate_cluster_name(name):
|
||||
continue
|
||||
path = sub / JOURNAL_SUBDIR / JOBS_HISTORY_FILENAME
|
||||
if not path.is_file():
|
||||
continue
|
||||
try:
|
||||
data = json.loads(path.read_text(encoding="utf-8"))
|
||||
except (OSError, json.JSONDecodeError):
|
||||
continue
|
||||
entries = data.get("entries") if isinstance(data, dict) else None
|
||||
if not isinstance(entries, list):
|
||||
continue
|
||||
for e in entries:
|
||||
if isinstance(e, dict):
|
||||
flat.append((name, e))
|
||||
|
||||
def sort_key(item: tuple[str, dict[str, Any]]) -> str:
|
||||
return str(item[1].get("finished_at_utc") or item[1].get("created_at_utc") or "")
|
||||
|
||||
flat.sort(key=sort_key, reverse=True)
|
||||
return flat
|
||||
|
||||
|
||||
def collect_recent_journal_entries_page_sync(*, limit: int, offset: int) -> tuple[list[dict[str, Any]], int]:
|
||||
"""
|
||||
Страница записей из всех ``clusters/*/journal/jobs_history.json``, новые первыми.
|
||||
|
||||
В каждую запись добавляется поле ``source_cluster`` (имя каталога).
|
||||
Возвращает ``(страница записей, общее число записей)``.
|
||||
"""
|
||||
lim = max(1, min(limit, 100))
|
||||
off = max(0, offset)
|
||||
flat = _collect_all_journal_rows_sorted_sync()
|
||||
total = len(flat)
|
||||
chunk = flat[off : off + lim]
|
||||
out: list[dict[str, Any]] = []
|
||||
for cluster_dir_name, e in chunk:
|
||||
row = dict(e)
|
||||
row["source_cluster"] = cluster_dir_name
|
||||
out.append(row)
|
||||
return out, total
|
||||
@@ -386,34 +386,87 @@ class JobStore:
|
||||
|
||||
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)
|
||||
snapshot: dict[str, Any] | None = None
|
||||
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
|
||||
r = self._jobs[job_id]
|
||||
snapshot = {
|
||||
"job_id": r.job_id,
|
||||
"kind": r.kind,
|
||||
"cluster_name": r.cluster_name,
|
||||
"status": r.status,
|
||||
"message": r.message,
|
||||
"created_at_utc": r.created_at_utc,
|
||||
"log_lines": list(r.log_lines),
|
||||
"result": dict(r.result) if r.result else None,
|
||||
}
|
||||
set_progress_sync(job_id, "Готово", 100)
|
||||
await self._persist_async()
|
||||
if snapshot and snapshot.get("cluster_name"):
|
||||
from core.job_journal import append_job_from_record_sync
|
||||
|
||||
await asyncio.to_thread(append_job_from_record_sync, **snapshot)
|
||||
|
||||
async def set_failed(self, job_id: str, message: str) -> None:
|
||||
logs = take_logs_finalize_sync(job_id)
|
||||
# Чтобы в kind_k8s_jobs.json и в UI всегда был хотя бы текст ошибки (ранний сбой без stdout).
|
||||
if not logs:
|
||||
logs = [f"[ошибка] {message}"]
|
||||
snapshot: dict[str, Any] | None = None
|
||||
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
|
||||
r = self._jobs[job_id]
|
||||
snapshot = {
|
||||
"job_id": r.job_id,
|
||||
"kind": r.kind,
|
||||
"cluster_name": r.cluster_name,
|
||||
"status": r.status,
|
||||
"message": r.message,
|
||||
"created_at_utc": r.created_at_utc,
|
||||
"log_lines": list(r.log_lines),
|
||||
"result": dict(r.result) if r.result else None,
|
||||
}
|
||||
logger.warning("Задание %s завершилось ошибкой: %s", job_id, message)
|
||||
await self._persist_async()
|
||||
if snapshot and snapshot.get("cluster_name"):
|
||||
from core.job_journal import append_job_from_record_sync
|
||||
|
||||
await asyncio.to_thread(append_job_from_record_sync, **snapshot)
|
||||
|
||||
async def set_cancelled(self, job_id: str, message: str = "Операция отменена") -> None:
|
||||
logs = take_logs_finalize_sync(job_id)
|
||||
if not logs:
|
||||
logs = [f"[отмена] {message}"]
|
||||
snapshot: dict[str, Any] | None = None
|
||||
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
|
||||
r = self._jobs[job_id]
|
||||
snapshot = {
|
||||
"job_id": r.job_id,
|
||||
"kind": r.kind,
|
||||
"cluster_name": r.cluster_name,
|
||||
"status": r.status,
|
||||
"message": r.message,
|
||||
"created_at_utc": r.created_at_utc,
|
||||
"log_lines": list(r.log_lines),
|
||||
"result": dict(r.result) if r.result else None,
|
||||
}
|
||||
logger.info("Задание %s отменено: %s", job_id, message)
|
||||
await self._persist_async()
|
||||
if snapshot and snapshot.get("cluster_name"):
|
||||
from core.job_journal import append_job_from_record_sync
|
||||
|
||||
await asyncio.to_thread(append_job_from_record_sync, **snapshot)
|
||||
|
||||
async def get(self, job_id: str) -> JobRecord | None:
|
||||
async with self._lock:
|
||||
|
||||
Reference in New Issue
Block a user