Панель: журнал заданий, pull --progress plain, PTY/EIO, старт/стоп в фоне, очистка jobs

- Фоновые stop/start с job_id и poll; отмена с kill; docker pull plain + снятие ANSI
- Лимиты журнала API/буфера; список jobs без progress_log; DELETE /jobs
- UI: опрос чаще, подсказка при пустом логе, кнопка очистки завершённых
This commit is contained in:
Sergey Antropoff
2026-04-04 07:04:46 +03:00
parent 6d4bc65c8a
commit 8bd44adbb0
10 changed files with 808 additions and 200 deletions
+348 -64
View File
@@ -8,6 +8,8 @@
from __future__ import annotations
import codecs
import errno
import json
import logging
import os
@@ -25,6 +27,18 @@ from kubeconfig_patch import patch_kubeconfig_server_for_host, should_patch_afte
logger = logging.getLogger("kind_k8s.cluster_lifecycle")
# Удаление ANSI/OSC из потока (docker TTY: курсор [1A, очистка [2K и т.д.) — в журнале остаётся читаемый текст.
_ANSI_CSI_RE = re.compile(r"\x1b\[[\d;?]*[ -/]*[@-~]")
_ANSI_OSC_RE = re.compile(r"\x1b\][^\x07]*(?:\x07|\x1b\\)")
def _strip_ansi_stream_text(s: str) -> str:
"""Убрать CSI/OSC-последовательности терминала из строки лога (docker TTY: [1A, [2K, …)."""
if not s:
return ""
t = _ANSI_CSI_RE.sub("", s)
return _ANSI_OSC_RE.sub("", t)
def _rollback_after_cancel(*, cluster_name: str, out_dir: Path) -> None:
"""Удалить кластер kind и каталог данных после запроса отмены (best-effort)."""
@@ -107,40 +121,222 @@ def _run_checked(cmd: list[str], *, cwd: Path | None = None) -> None:
raise KindClusterError(f"Команда завершилась с кодом {p.returncode}: {err}", exit_code=p.returncode)
def _stream_pty_enabled() -> bool:
"""
Псевдо-TTY: docker/kind при выводе в pipe часто полностью буферизуют stdout — в журнале «тишина»
до конца команды. С TTY строки прогресса (в т.ч. pull) уходят порциями.
Отключение: ``KIND_K8S_STREAM_PTY=0``.
"""
if os.name != "posix":
return False
v = (os.environ.get("KIND_K8S_STREAM_PTY") or "1").strip().lower()
return v not in ("0", "false", "no", "off", "нет")
def _emit_stream_lines(
buffer: str,
on_line: Callable[[str], None] | None,
*,
flush_tail: bool,
) -> str:
"""
Разобрать буфер по ``\\n`` и ``\\r`` (прогресс docker/containerd часто идёт с ``\\r`` без ``\\n``).
При ``flush_tail=False`` возвращает неполный хвост после последнего ``\\n``.
"""
if not buffer:
return ""
normalized = buffer.replace("\r\n", "\n").replace("\r", "\n")
parts = normalized.split("\n")
if flush_tail:
for part in parts:
s = _strip_ansi_stream_text(part).strip()
if s and on_line:
on_line(s)
if s:
logger.debug("stream: %s", s[:800])
return ""
if len(parts) == 1:
return parts[0]
*body, tail = parts
for part in body:
s = _strip_ansi_stream_text(part).strip()
if s and on_line:
on_line(s)
if s:
logger.debug("stream: %s", s[:800])
return tail
def _run_checked_stream(
cmd: list[str],
*,
cwd: Path | None = None,
on_line: Callable[[str], None] | None = None,
job_id: str | None = None,
use_pty: bool | None = None,
) -> None:
"""
Выполнить команду с построчным выводом в колбэк (stdout+stderr объединены).
Выполнить команду с потоковым выводом в колбэк (stdout+stderr объединены).
Нужен для ``kind create cluster``: pull образов и подъём нод видны в UI по опросу job.
``use_pty=None`` — как в ``KIND_K8S_STREAM_PTY``. ``use_pty=False`` — только pipe (удобно вместе
с ``docker pull --progress plain``). ``use_pty=True`` — принудительно PTY на POSIX.
``job_id`` — регистрация ``Popen`` для принудительной отмены (``SIGKILL`` группы процессов).
"""
from core import job_store as _js
logger.info("Выполнение (поток): %s", " ".join(cmd))
p = subprocess.Popen(
cmd,
cwd=cwd,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
bufsize=1,
)
if p.stdout is None:
raise KindClusterError("Не удалось открыть stdout процесса", exit_code=1)
master_fd: int | None = None
if use_pty is None:
use_pty = _stream_pty_enabled()
slave_fd: int | None = None
if use_pty:
try:
import pty
master_fd, slave_fd = pty.openpty()
except OSError as e:
logger.warning("PTY недоступен (%s), поток через pipe — вывод может обновляться с задержкой", e)
use_pty = False
master_fd = None
slave_fd = None
if use_pty and slave_fd is not None and master_fd is not None:
try:
p = subprocess.Popen(
cmd,
cwd=cwd,
stdin=slave_fd,
stdout=slave_fd,
stderr=subprocess.STDOUT,
close_fds=False,
start_new_session=True,
)
except Exception:
try:
os.close(slave_fd)
except OSError:
pass
try:
os.close(master_fd)
except OSError:
pass
raise
try:
os.close(slave_fd)
except OSError:
pass
slave_fd = None
else:
p = subprocess.Popen(
cmd,
cwd=cwd,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
bufsize=0,
start_new_session=True,
)
master_fd = None
if p.stdout is None:
raise KindClusterError("Не удалось открыть stdout процесса", exit_code=1)
if job_id:
_js.register_subprocess_sync(job_id, p)
decoder = codecs.getincrementaldecoder("utf-8")(errors="replace")
carry = ""
rc = -1
try:
for raw in p.stdout:
line = raw.rstrip("\n\r")
if on_line and line:
on_line(line)
if line:
logger.debug("stream: %s", line[:800])
while True:
try:
if master_fd is not None:
raw = os.read(master_fd, 4096)
else:
raw = p.stdout.read(4096) if p.stdout else b""
except OSError as e:
# EINTR — повторить read; EIO — slave PTY закрыт (часто после SIGKILL при отмене).
_eio = getattr(errno, "EIO", 5)
if e.errno == errno.EINTR:
continue
if e.errno == _eio:
logger.debug("Чтение потока: EIO (вероятно процесс завершён), выходим из цикла")
break
raise
if not raw:
break
carry += decoder.decode(raw)
carry = _emit_stream_lines(carry, on_line, flush_tail=False)
carry += decoder.decode(b"", final=True)
carry = _emit_stream_lines(carry, on_line, flush_tail=True)
rc = p.wait()
finally:
p.stdout.close()
if job_id:
_js.unregister_subprocess_sync(job_id)
if master_fd is not None:
try:
os.close(master_fd)
except OSError:
pass
else:
try:
if p.stdout:
p.stdout.close()
except OSError:
pass
if job_id and _js.is_cancelled_sync(job_id):
raise KindClusterError("Операция отменена", exit_code=130)
if rc != 0:
raise KindClusterError(f"Команда завершилась с кодом {rc} (см. журнал задания выше)", exit_code=rc)
raise KindClusterError(f"Команда завершилась с кодом {rc} (см. журнал выше)", exit_code=rc)
def _docker_pull_plain_progress_enabled() -> bool:
"""
``docker pull --progress plain`` — построчный текст без TTY-ANSI; удобно с pipe (без PTY).
Отключить: ``KIND_K8S_DOCKER_PULL_PLAIN=0`` (если старая версия docker не знает флаг).
"""
v = (os.environ.get("KIND_K8S_DOCKER_PULL_PLAIN") or "1").strip().lower()
return v not in ("0", "false", "no", "off", "нет")
def _pull_node_image_streaming(*, image: str, on_line: Callable[[str], None], job_id: str | None) -> None:
"""
Скачать образ узлов через ``docker pull`` / ``podman pull`` с выводом в журнал задания.
Для **docker** по умолчанию: ``--progress plain`` и **без PTY** — идут строки Downloading/Extracting
с процентами. Для **podman** — обычный pull и PTY по ``KIND_K8S_STREAM_PTY``.
"""
cli = _container_cli_bin()
if not shutil.which(cli):
raise KindClusterError(
"Скачивание образа недоступно: не найдена программа для работы с контейнерами (docker/podman).",
exit_code=127,
)
if cli == "docker" and _docker_pull_plain_progress_enabled():
try:
_run_checked_stream(
[cli, "pull", "--progress=plain", image],
on_line=on_line,
job_id=job_id,
use_pty=False,
)
return
except KindClusterError as e:
# У docker CLI код 125 — «неверная команда/флаг»; тогда пробуем pull без plain.
err_low = str(e).lower()
if e.exit_code == 125 or any(
x in err_low
for x in ("progress", "unknown flag", "bad flag", "invalid option", "unknown shorthand")
):
logger.warning("Повтор docker pull без --progress=plain: %s", e)
_run_checked_stream([cli, "pull", image], on_line=on_line, job_id=job_id, use_pty=None)
return
raise
_run_checked_stream([cli, "pull", image], on_line=on_line, job_id=job_id, use_pty=None)
def _run_capture_checked(cmd: list[str]) -> str:
@@ -231,7 +427,6 @@ def create_cluster_non_interactive(
def _progress(stage: str, pct: int) -> None:
if job_id:
_job_store.set_progress_sync(job_id, stage, pct)
_job_store.append_log_sync(job_id, f"[{pct}%] {stage}")
def _log(line: str) -> None:
if job_id:
@@ -274,15 +469,17 @@ def create_cluster_non_interactive(
node_image = str(prev["node_image"])
if prev.get("kubernetes_version_tag"):
ver_tag = str(prev["kubernetes_version_tag"])
_progress("Используется существующий kind-config.yaml", 10)
_progress("Чтение сохранённого конфига", 10)
_log("Используется сохранённый файл конфигурации кластера.")
else:
yaml_text = build_kind_config_yaml(node_image=node_image, workers=workers)
cfg_path.write_text(yaml_text, encoding="utf-8")
_progress("Подготовка каталога и kind-config", 12)
_progress("Подготовка конфигурации", 12)
_log("Файл конфигурации кластера сохранён в каталог данных.")
if _cancelled():
_rollback_after_cancel(cluster_name=name, out_dir=out_dir)
raise KindClusterError("Создание отменено пользователем")
raise KindClusterError("Операция отменена")
logger.info(
"Создание кластера «%s», образ %s, workers=%s, existing_cfg=%s",
@@ -291,38 +488,54 @@ def create_cluster_non_interactive(
workers,
use_existing_config,
)
_progress("kind create cluster (скачивание образов и подъём нод — может занять несколько минут)", 28)
_log("--- kind create cluster ---")
_progress("Скачивание образа при необходимости", 18)
if _cancelled():
_rollback_after_cancel(cluster_name=name, out_dir=out_dir)
raise KindClusterError("Операция отменена")
try:
_pull_node_image_streaming(image=node_image, on_line=_log, job_id=job_id)
except KindClusterError as e:
if _cancelled():
_rollback_after_cancel(cluster_name=name, out_dir=out_dir)
raise KindClusterError("Операция отменена") from e
raise
if _cancelled():
_rollback_after_cancel(cluster_name=name, out_dir=out_dir)
raise KindClusterError("Операция отменена")
_progress("Создание узлов кластера", 32)
_run_checked_stream(
["kind", "create", "cluster", "--name", name, "--config", str(cfg_path)],
on_line=_log,
job_id=job_id,
)
if _cancelled():
_rollback_after_cancel(cluster_name=name, out_dir=out_dir)
raise KindClusterError("Создание отменено пользователем")
raise KindClusterError("Операция отменена")
_progress("Сохранение kubeconfig", 58)
_progress("Сохранение доступа к кластеру", 58)
kube = _run_capture_checked(["kind", "get", "kubeconfig", "--name", name])
kube_path.write_text(kube, encoding="utf-8")
if _cancelled():
_rollback_after_cancel(cluster_name=name, out_dir=out_dir)
raise KindClusterError("Создание отменено пользователем")
raise KindClusterError("Операция отменена")
patched = False
if should_patch_after_create():
_progress("Патч kubeconfig для доступа к API с хоста", 72)
_progress("Настройка доступа с вашего компьютера", 72)
patched = patch_kubeconfig_server_for_host(cluster_name=name, kube_path=kube_path)
if _cancelled():
_rollback_after_cancel(cluster_name=name, out_dir=out_dir)
raise KindClusterError("Создание отменено пользователем")
raise KindClusterError("Операция отменена")
nodes_ready: bool | None = None
nodes_msg: str | None = None
if _wait_nodes_enabled():
_progress("Ожидание готовности нод (kubectl wait …)", 82)
_progress("Ожидание готовности узлов", 82)
ok, msg = wait_nodes_ready(kubeconfig_path=kube_path)
nodes_ready = ok
nodes_msg = msg
@@ -330,7 +543,7 @@ def create_cluster_non_interactive(
logger.info("Ноды готовы: %s", msg)
else:
logger.warning("Ожидание нод не завершилось успешно: %s", msg)
_log(f"kubectl wait nodes: {msg}"[:4000])
_log(("Узлы готовы" if ok else "Ожидание узлов: ") + (msg or "")[:4000])
worker_nodes_meta = workers
if use_existing_config:
@@ -357,7 +570,7 @@ def create_cluster_non_interactive(
}
meta_path.write_text(json.dumps(meta, ensure_ascii=False, indent=2), encoding="utf-8")
_progress("Финализация", 95)
_progress("Завершение", 95)
return CreateClusterResult(
cluster_name=name,
@@ -434,7 +647,7 @@ def list_kind_cluster_container_names(*, cluster_name: str) -> list[str]:
"""Имена контейнеров узлов kind (все с префиксом ``<имя>-``)."""
cli = _container_cli_bin()
if not shutil.which(cli):
raise KindClusterError(f"CLI контейнеров «{cli}» не найден в PATH.", exit_code=127)
raise KindClusterError("Программа для контейнеров не найдена в PATH.", exit_code=127)
p = subprocess.run(
[cli, "ps", "-a", "--format", "{{.Names}}"],
capture_output=True,
@@ -442,59 +655,130 @@ def list_kind_cluster_container_names(*, cluster_name: str) -> list[str]:
)
if p.returncode != 0:
err = (p.stderr or p.stdout or "").strip()
raise KindClusterError(f"{cli} ps: {err}", exit_code=p.returncode)
raise KindClusterError(f"Не удалось получить список контейнеров: {err}", exit_code=p.returncode)
prefix = f"{cluster_name}-"
raw = [n.strip() for n in (p.stdout or "").splitlines() if n.strip()]
matched = [n for n in raw if n.startswith(prefix)]
return _sort_kind_node_containers(matched)
def stop_kind_cluster_containers(*, name: str) -> tuple[bool, str]:
def stop_kind_cluster_containers(*, name: str, job_id: str | None = None) -> tuple[bool, str]:
"""
Остановить контейнеры узлов (``docker stop`` / ``podman stop``).
Остановить контейнеры узлов (CLI контейнеров из ``CONTAINER_CLI``).
Запись kind о кластере сохраняется; позже можно вызвать ``start_kind_cluster_containers``.
При ``job_id`` — потоковый вывод в журнал задания и возможность прервать текущую команду отменой.
"""
from core import job_store as _js
names = list_kind_cluster_container_names(cluster_name=name)
if not names:
return True, "Нет контейнеров с префиксом «%s-» (уже остановлены или удалены)" % name
return True, "Узлы уже остановлены или отсутствуют"
def _cancelled() -> bool:
return bool(job_id and _js.is_cancelled_sync(job_id))
cli = _container_cli_bin()
n = len(names)
ok_all = True
parts: list[str] = []
for ctr in names:
p = subprocess.run([cli, "stop", ctr], capture_output=True, text=True)
if p.returncode != 0:
ok_all = False
err = (p.stderr or p.stdout or "").strip() or str(p.returncode)
parts.append(f"{ctr}: ошибка ({err})")
logger.warning("%s stop %s: %s", cli, ctr, err)
failed_idx: list[int] = []
for i, ctr in enumerate(names, start=1):
if _cancelled():
raise KindClusterError("Операция отменена", exit_code=130)
if job_id:
pct = 5 + int(88 * (i - 1) / max(n, 1))
_js.set_progress_sync(job_id, f"Остановка узла {i} из {n}", min(pct, 93))
_js.append_log_sync(job_id, f"Узел {i} из {n}")
try:
def _on_line(line: str) -> None:
_js.append_log_sync(job_id, line)
_run_checked_stream([cli, "stop", ctr], on_line=_on_line, job_id=job_id)
except KindClusterError as e:
if _cancelled():
raise KindClusterError("Операция отменена", exit_code=130) from e
ok_all = False
failed_idx.append(i)
logger.warning("%s stop узел %s: %s", cli, i, e)
else:
parts.append(f"{ctr}: OK")
return ok_all, "; ".join(parts)
p = subprocess.run([cli, "stop", ctr], capture_output=True, text=True)
if p.returncode != 0:
ok_all = False
err = (p.stderr or p.stdout or "").strip() or str(p.returncode)
failed_idx.append(i)
logger.warning("%s stop %s: %s", cli, ctr, err)
if job_id:
_js.set_progress_sync(job_id, "Готово", 98)
if ok_all:
return True, f"Остановлено узлов: {n}"
return False, "Не удалось остановить узлы: " + (
", ".join(str(x) for x in failed_idx) if failed_idx else "см. журнал"
)
def start_kind_cluster_containers(*, name: str) -> tuple[bool, str]:
"""Запустить контейнеры узлов kind (после ``stop`` или рестарта движка)."""
def start_kind_cluster_containers(*, name: str, job_id: str | None = None) -> tuple[bool, str]:
"""
Запустить контейнеры узлов (после остановки или перезапуска движка контейнеров).
При ``job_id`` — журнал и прогресс в задании, прерывание текущей команды через отмену задания.
"""
from core import job_store as _js
names = list_kind_cluster_container_names(cluster_name=name)
if not names:
return False, (
"Не найдены контейнеры «%s-*». Если кластера нет в kind — используйте «Старт» "
"из UI (создание по сохранённому kind-config.yaml) или создайте кластер заново."
% name
"Сохранённые узлы не найдены. Если кластер ещё не создавался в этой среде — создайте его "
"или поднимите по сохранённой конфигурации из панели."
)
def _cancelled() -> bool:
return bool(job_id and _js.is_cancelled_sync(job_id))
cli = _container_cli_bin()
n = len(names)
ok_all = True
parts: list[str] = []
for ctr in names:
p = subprocess.run([cli, "start", ctr], capture_output=True, text=True)
if p.returncode != 0:
ok_all = False
err = (p.stderr or p.stdout or "").strip() or str(p.returncode)
parts.append(f"{ctr}: ошибка ({err})")
logger.warning("%s start %s: %s", cli, ctr, err)
failed_idx: list[int] = []
for i, ctr in enumerate(names, start=1):
if _cancelled():
raise KindClusterError("Операция отменена", exit_code=130)
if job_id:
pct = 5 + int(88 * (i - 1) / max(n, 1))
_js.set_progress_sync(job_id, f"Запуск узла {i} из {n}", min(pct, 93))
_js.append_log_sync(job_id, f"Узел {i} из {n}")
try:
def _on_line(line: str) -> None:
_js.append_log_sync(job_id, line)
_run_checked_stream([cli, "start", ctr], on_line=_on_line, job_id=job_id)
except KindClusterError as e:
if _cancelled():
raise KindClusterError("Операция отменена", exit_code=130) from e
ok_all = False
failed_idx.append(i)
logger.warning("%s start узел %s: %s", cli, i, e)
else:
parts.append(f"{ctr}: OK")
return ok_all, "; ".join(parts)
p = subprocess.run([cli, "start", ctr], capture_output=True, text=True)
if p.returncode != 0:
ok_all = False
err = (p.stderr or p.stdout or "").strip() or str(p.returncode)
failed_idx.append(i)
logger.warning("%s start %s: %s", cli, ctr, err)
if job_id:
_js.set_progress_sync(job_id, "Готово", 98)
if ok_all:
return True, f"Запущено узлов: {n}"
return False, "Не удалось запустить узлы: " + (
", ".join(str(x) for x in failed_idx) if failed_idx else "см. журнал"
)
def read_meta_json(cluster_name: str) -> dict[str, object] | None: