From f5a808873a905bc1f220d11837aa7562b5168830 Mon Sep 17 00:00:00 2001 From: Sergey Antropoff Date: Sat, 4 Apr 2026 07:04:46 +0300 Subject: [PATCH] =?UTF-8?q?=D0=9F=D0=B0=D0=BD=D0=B5=D0=BB=D1=8C:=20=D0=B6?= =?UTF-8?q?=D1=83=D1=80=D0=BD=D0=B0=D0=BB=20=D0=B7=D0=B0=D0=B4=D0=B0=D0=BD?= =?UTF-8?q?=D0=B8=D0=B9,=20pull=20--progress=20plain,=20PTY/EIO,=20=D1=81?= =?UTF-8?q?=D1=82=D0=B0=D1=80=D1=82/=D1=81=D1=82=D0=BE=D0=BF=20=D0=B2=20?= =?UTF-8?q?=D1=84=D0=BE=D0=BD=D0=B5,=20=D0=BE=D1=87=D0=B8=D1=81=D1=82?= =?UTF-8?q?=D0=BA=D0=B0=20jobs?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Фоновые stop/start с job_id и poll; отмена с kill; docker pull plain + снятие ANSI - Лимиты журнала API/буфера; список jobs без progress_log; DELETE /jobs - UI: опрос чаще, подсказка при пустом логе, кнопка очистки завершённых --- README.md | 5 +- app/api/v1/endpoints/clusters.py | 192 ++++++++++---- app/core/cluster_lifecycle.py | 412 ++++++++++++++++++++++++++----- app/core/job_store.py | 79 +++++- app/docs/api_routes.md | 88 ++++--- app/models/schemas.py | 4 +- app/static/js/dashboard.js | 196 ++++++++++++--- app/static/style.css | 18 +- app/templates/dashboard.html | 9 +- docker-compose.yml | 5 + 10 files changed, 808 insertions(+), 200 deletions(-) diff --git a/README.md b/README.md index cb59e78..3bb1dda 100644 --- a/README.md +++ b/README.md @@ -150,7 +150,10 @@ docker compose exec kind-k8s-web kubectl --kubeconfig=/work/clusters/<имя>/ku | **`KIND_K8S_VERSION_LIST_DISPLAY`** | контейнер | Сколько тегов отдавать в API/UI | | **`KIND_K8S_HUB_TAGS_MAX_PAGES`** | контейнер | Лимит страниц API Hub | | **`KIND_K8S_DEBUG`** | контейнер | `1`/`true`/`yes`/`да` — уровень DEBUG в логах | -| **`KIND_K8S_JOB_LOG_MAX_LINES`** | приложение | Размер буфера строк журнала фонового задания (`kind create`) для поля `progress_log` в API/UI; по умолчанию **500** (задаётся в коде, при необходимости передайте в compose) | +| **`KIND_K8S_JOB_LOG_MAX_LINES`** | приложение | Сколько строк журнала хранить в памяти на задание (старые вытесняются); по умолчанию **2500** | +| **`KIND_K8S_JOB_API_LOG_MAX_LINES`** | приложение | Сколько строк отдавать в **GET /api/v1/jobs/{id}** (хвост); по умолчанию **5000**, максимум **20000** | +| **`KIND_K8S_STREAM_PTY`** | приложение | **`1`** (по умолчанию) — для `kind` и **podman pull** псевдо-TTY; **`0`** — только pipe | +| **`KIND_K8S_DOCKER_PULL_PLAIN`** | приложение | **`1`** (по умолчанию) — `docker pull --progress=plain` без PTY (построчный прогресс в журнале); **`0`** — обычный pull | | **`KIND_K8S_README_PATH`** | контейнер / приложение | Абсолютный путь к **README.md** для страницы **`/documentation`**; если пусто — используется `README.md` рядом с каталогом `app/` (в образе: `/opt/kind-k8s/README.md`) | | **`KIND_K8S_WORKDIR`** | локальный запуск | Корень данных на машине разработчика без compose | | **`COMPOSE_BUILD_FLAGS`** | Makefile | Например `make docker build COMPOSE_BUILD_FLAGS=--platform linux/arm64` (то же для **`make docker rebuild`**) | diff --git a/app/api/v1/endpoints/clusters.py b/app/api/v1/endpoints/clusters.py index 59cfc24..87deab5 100644 --- a/app/api/v1/endpoints/clusters.py +++ b/app/api/v1/endpoints/clusters.py @@ -8,6 +8,7 @@ from __future__ import annotations import asyncio import logging +import os from typing import Any from fastapi import APIRouter, BackgroundTasks, HTTPException, Query @@ -51,19 +52,31 @@ logger = logging.getLogger("kind_k8s.api.clusters") router = APIRouter(tags=["clusters"]) -def _record_to_job_view(rec: JobRecord) -> JobView: +def _api_job_log_limit_detail() -> int: + """Сколько строк журнала отдавать в GET /jobs/{id} (полный хвост буфера задания).""" + raw = (os.environ.get("KIND_K8S_JOB_API_LOG_MAX_LINES") or "5000").strip() + try: + return max(200, min(int(raw), 20000)) + except ValueError: + return 5000 + + +def _record_to_job_view(rec: JobRecord, *, include_progress_log: bool = True) -> JobView: """JobRecord → JobView с полями прогресса из потокобезопасного снимка.""" prog = get_progress_sync(rec.job_id) stage, pct = (None, None) if prog is not None: stage, pct = prog[0], prog[1] - if rec.status in ("queued", "running"): - log_tail = get_logs_snapshot_sync(rec.job_id) + if not include_progress_log: + log_tail: list[str] = [] else: - log_tail = list(rec.log_lines or []) - max_log = 400 - if len(log_tail) > max_log: - log_tail = log_tail[-max_log:] + if rec.status in ("queued", "running"): + log_tail = get_logs_snapshot_sync(rec.job_id) + else: + log_tail = list(rec.log_lines or []) + cap = _api_job_log_limit_detail() + if len(log_tail) > cap: + log_tail = log_tail[-cap:] return JobView( job_id=rec.job_id, kind=rec.kind, @@ -127,7 +140,18 @@ async def get_stats() -> StatsResponse: async def list_jobs(limit: int = Query(30, ge=1, le=200, description="Сколько последних заданий")) -> list[JobView]: """История создания кластеров (в памяти процесса; после перезапуска контейнера пусто).""" items = job_store.snapshot_recent_sorted(limit=limit) - return [_record_to_job_view(r) for r in items] + # В списке не тащим progress_log — экономим трафик; полный журнал только в GET /jobs/{id}. + return [_record_to_job_view(r, include_progress_log=False) for r in items] + + +@router.delete("/jobs", summary="Очистить завершённые задания из памяти") +async def delete_completed_jobs() -> dict[str, int]: + """ + Удаляет записи со статусом success, failed, cancelled. Задания в очереди и в работе не трогаются. + """ + removed = await job_store.purge_completed() + logger.info("Очистка заданий: удалено завершённых записей: %s", removed) + return {"removed": removed} @router.get("/clusters", response_model=list[ClusterSummary], summary="Список кластеров") @@ -230,7 +254,7 @@ async def _run_create_job(job_id: str, body: ClusterCreateRequest) -> None: ) except KindClusterError as e: msg = str(e) - if "отменено" in msg.lower(): + if "отмен" in msg.lower(): await job_store.set_cancelled(job_id, msg) else: await job_store.set_failed(job_id, msg) @@ -273,7 +297,7 @@ async def _run_start_cluster_job(job_id: str, name: str, kubernetes_version_tag: ) except KindClusterError as e: msg = str(e) - if "отменено" in msg.lower(): + if "отмен" in msg.lower(): await job_store.set_cancelled(job_id, msg) else: await job_store.set_failed(job_id, msg) @@ -300,6 +324,72 @@ async def _run_start_cluster_job(job_id: str, name: str, kubernetes_version_tag: end_job_tracking(job_id) +async def _run_stop_cluster_job(job_id: str, name: str) -> None: + """Фоновая остановка узлов с журналом в GET /jobs/{id}.""" + n = name.strip() + try: + async with kind_cluster_lock: + await job_store.set_running(job_id, stage="Остановка узлов", percent=6) + try: + ok, summary = await asyncio.to_thread(stop_kind_cluster_containers, name=n, job_id=job_id) + except KindClusterError as e: + msg = str(e) + if "отмен" in msg.lower(): + await job_store.set_cancelled(job_id, msg) + else: + await job_store.set_failed(job_id, msg) + logger.warning("stop_containers job %s: %s", job_id, e) + return + except Exception as e: + await job_store.set_failed(job_id, f"{type(e).__name__}: {e}") + logger.exception("stop_containers job %s: непредвиденная ошибка", job_id) + return + + await job_store.set_success( + job_id, + result={"name": n, "containers_stopped_ok": ok}, + message=summary, + ) + logger.info("stop_containers job %s: ok=%s", job_id, ok) + finally: + end_job_tracking(job_id) + + +async def _run_start_containers_job(job_id: str, name: str) -> None: + """Фоновый запуск уже существующих узлов (кластер зарегистрирован в kind).""" + n = name.strip() + try: + async with kind_cluster_lock: + await job_store.set_running(job_id, stage="Запуск узлов", percent=6) + try: + ok, summary = await asyncio.to_thread(start_kind_cluster_containers, name=n, job_id=job_id) + except KindClusterError as e: + msg = str(e) + if "отмен" in msg.lower(): + await job_store.set_cancelled(job_id, msg) + else: + await job_store.set_failed(job_id, msg) + logger.warning("start_containers job %s: %s", job_id, e) + return + except Exception as e: + await job_store.set_failed(job_id, f"{type(e).__name__}: {e}") + logger.exception("start_containers job %s: непредвиденная ошибка", job_id) + return + + if not ok: + await job_store.set_failed(job_id, summary) + return + + await job_store.set_success( + job_id, + result={"name": n, "containers_started_ok": ok}, + message=summary, + ) + logger.info("start_containers job %s: ok=%s", job_id, ok) + finally: + end_job_tracking(job_id) + + @router.post( "/clusters", response_model=ClusterCreateAccepted, @@ -346,30 +436,31 @@ async def delete_cluster(name: str) -> dict[str, object]: @router.post( "/clusters/{name}/stop", - summary="Остановить узлы кластера (docker stop)", + summary="Остановить узлы кластера (фон, журнал в /jobs)", responses={400: {"description": "Некорректное имя"}}, ) -async def stop_cluster_nodes(name: str) -> dict[str, object]: +async def stop_cluster_nodes(name: str, background_tasks: BackgroundTasks) -> JSONResponse: """ - Остановить контейнеры узлов kind; запись кластера в kind сохраняется. + Поставить остановку узлов в фон: тот же журнал в реальном времени, что и при создании. - После этого API «Старт» запустит те же контейнеры без ``kind create``. + Запись кластера в kind сохраняется; «Старт» снова поднимет узлы. """ if not validate_cluster_name(name): raise HTTPException(status_code=400, detail="Некорректное имя кластера") - async with kind_cluster_lock: - - def _do() -> tuple[bool, str]: - return stop_kind_cluster_containers(name=name) - - try: - ok, summary = await asyncio.to_thread(_do) - except KindClusterError as e: - raise HTTPException(status_code=500, detail=str(e)) from e - - logger.info("Остановка узлов %s: ok=%s", name, ok) - return {"name": name, "containers_stopped_ok": ok, "summary": summary} + n = name.strip() + rec = await job_store.create_job("stop_containers", cluster_name=n) + background_tasks.add_task(_run_stop_cluster_job, rec.job_id, n) + logger.info("Остановка узлов %s в фоне, job_id=%s", n, rec.job_id) + return JSONResponse( + status_code=202, + content={ + "job_id": rec.job_id, + "status": "queued", + "mode": "stop", + "message": "Остановка узлов; опросите GET /api/v1/jobs/{job_id}", + }, + ) @router.post( @@ -382,10 +473,9 @@ async def start_cluster_nodes( background_tasks: BackgroundTasks, ) -> JSONResponse: """ - Если кластер есть в ``kind get clusters`` — ``docker start`` всех узлов. + Если кластер уже зарегистрирован — фоновый запуск сохранённых узлов (журнал в GET /jobs). - Если в kind нет, но есть ``clusters/<имя>/kind-config.yaml`` — фоновое ``kind create`` - (как при создании, с журналом в GET /jobs/{job_id}). + Иначе, при наличии сохранённого конфига в каталоге кластера — фоновый подъём как при создании. """ if not validate_cluster_name(name): raise HTTPException(status_code=400, detail="Некорректное имя кластера") @@ -394,25 +484,20 @@ async def start_cluster_nodes( async with kind_cluster_lock: in_kind = n in await asyncio.to_thread(list_registered_kind_clusters) - if in_kind: - def _start() -> tuple[bool, str]: - return start_kind_cluster_containers(name=n) - - try: - ok, summary = await asyncio.to_thread(_start) - except KindClusterError as e: - raise HTTPException(status_code=500, detail=str(e)) from e - logger.info("Запуск контейнеров кластера %s: ok=%s", n, ok) - return JSONResponse( - status_code=200, - content={ - "name": n, - "mode": "containers", - "containers_started_ok": ok, - "summary": summary, - }, - ) + if in_kind: + rec = await job_store.create_job("start_containers", cluster_name=n) + background_tasks.add_task(_run_start_containers_job, rec.job_id, n) + logger.info("Запуск контейнеров кластера %s в фоне, job_id=%s", n, rec.job_id) + return JSONResponse( + status_code=202, + content={ + "job_id": rec.job_id, + "status": "queued", + "mode": "containers", + "message": "Запуск узлов; опросите GET /api/v1/jobs/{job_id}", + }, + ) cfg = clusters_dir() / n / "kind-config.yaml" if not cfg.is_file(): @@ -437,22 +522,23 @@ async def start_cluster_nodes( content={ "job_id": rec.job_id, "status": "queued", - "message": "Подъём кластера по kind-config.yaml; опросите GET /api/v1/jobs/{job_id}", + "mode": "kind_config", + "message": "Подъём кластера по сохранённому конфигу; опросите GET /api/v1/jobs/{job_id}", }, ) @router.post( "/jobs/{job_id}/cancel", - summary="Запросить отмену создания кластера", + summary="Прервать фоновое задание", responses={400: {"description": "Задание уже завершено"}, 404: {"description": "Нет задания"}}, ) async def cancel_create_job(job_id: str) -> dict[str, object]: """ - Установить флаг отмены для задания ``create_cluster`` или ``start_cluster``. + Прервать активное задание: скачивание образа, создание кластера, запуск/остановка узлов. - Этап ``kind create cluster`` нельзя прервать до его завершения; после него отмена удалит - кластер и данные (если успели создать). + Для длительных команд дочерний процесс завершается принудительно; между узлами останов + также учитывает флаг отмены. """ rec = await job_store.get(job_id) if not rec: @@ -465,7 +551,7 @@ async def cancel_create_job(job_id: str) -> dict[str, object]: return { "job_id": job_id, "cancel_requested": True, - "message": "Отмена обрабатывается между этапами; во время kind create дождитесь окончания шага", + "message": "Запрошено прерывание; текущая команда будет остановлена, задание перейдёт в отменено", } diff --git a/app/core/cluster_lifecycle.py b/app/core/cluster_lifecycle.py index 8af0919..ccc2e15 100644 --- a/app/core/cluster_lifecycle.py +++ b/app/core/cluster_lifecycle.py @@ -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: diff --git a/app/core/job_store.py b/app/core/job_store.py index 9443ff5..c721710 100644 --- a/app/core/job_store.py +++ b/app/core/job_store.py @@ -13,6 +13,8 @@ from __future__ import annotations import asyncio import logging import os +import signal +import subprocess import threading import uuid from collections import deque @@ -33,10 +35,14 @@ _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 _max_job_log_lines() -> int: - raw = (os.environ.get("KIND_K8S_JOB_LOG_MAX_LINES") or "500").strip() + # Буфер в памяти: старые строки вытесняются; по умолчанию 2500 — виден почти весь типичный лог. + raw = (os.environ.get("KIND_K8S_JOB_LOG_MAX_LINES") or "2500").strip() try: return max(50, min(int(raw), 5000)) except ValueError: @@ -82,6 +88,8 @@ def begin_job_tracking(job_id: str) -> None: 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) @@ -102,19 +110,52 @@ def get_progress_sync(job_id: str) -> tuple[str, int] | None: 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: - """Запросить отмену. Вернуть False, если job_id не отслеживается.""" + """Запросить отмену и прервать активный подпроцесс (если есть).""" with _thread_lock: ev = _cancel_events.get(job_id) if ev is None: return False ev.set() logger.info("Запрошена отмена задания %s", job_id) - return True + kill_subprocess_for_job_sync(job_id) + return True def is_cancelled_sync(job_id: str) -> bool: - """Проверка из worker-thread между этапами создания кластера.""" + """Проверка из worker-thread между этапами и перед каждой следующей командой.""" with _thread_lock: ev = _cancel_events.get(job_id) return bool(ev and ev.is_set()) @@ -164,12 +205,18 @@ class JobStore: logger.info("Создано задание %s kind=%s cluster=%s", jid, kind, cluster_name) return rec - async def set_running(self, job_id: str) -> None: + 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, "Запуск создания кластера…", 5) + set_progress_sync(job_id, stage, percent) 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) @@ -190,7 +237,7 @@ class JobStore: self._jobs[job_id].log_lines = logs logger.warning("Задание %s завершилось ошибкой: %s", job_id, message) - async def set_cancelled(self, job_id: str, message: str = "Создание отменено пользователем") -> None: + 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: @@ -213,6 +260,24 @@ class JobStore: 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) + return len(to_remove) + # Синглтон на процесс uvicorn job_store = JobStore() diff --git a/app/docs/api_routes.md b/app/docs/api_routes.md index 99e259f..26e222a 100644 --- a/app/docs/api_routes.md +++ b/app/docs/api_routes.md @@ -19,7 +19,7 @@ | Маршрут | Описание | |---------|----------| -| `GET /` | HTML-панель: единая карточка «панель + среда», статистика, создание кластера (прогресс, **журнал** `kind create`, отмена), таблица кластеров с **иконками** действий и **всплывающими подсказками**, модалка узлов/подов; шапка — пилюли, Swagger / ReDoc / Health в отдельных окнах. | +| `GET /` | HTML-панель: единая карточка «панель + среда», статистика, создание кластера (прогресс, **журнал** операции в реальном времени, **отмена** с прерыванием текущей команды), таблица кластеров с **иконками** действий и **всплывающими подсказками**, модалка узлов/подов; шапка — пилюли, Swagger / ReDoc / Health в отдельных окнах. | | `GET /documentation` | HTML-оболочка; **`documentation.js`**: без `path` — **`GET /api/v1/docs/readme`**, с `?path=app/docs/…` — **`GET /api/v1/docs/file`**; разбор Markdown из **`/static/js/vendor/`** (marked, DOMPurify). Каждая секция по **H2** — **одна карточка** (заголовок h2 и содержимое до следующего h2 вместе). Заголовок вкладки браузера: **«Документация — …»** + текст **первого H1** документа + имя приложения (`KIND_K8S_APP_TITLE` на `body`). В шапке на этой странице активна только **Документация**; **Панель** как обычная пилюля (на дашборде активна **Панель**). Путь к README: `KIND_K8S_README_PATH` или `README.md` рядом с `app/`; в образе — `/opt/kind-k8s/README.md`. | | `GET /ui` | Редирект **307** на `/` (удобный ярлык). | | `GET /static/…` | CSS (`style.css`), скрипты панели (`js/dashboard.js`) и документации (`js/documentation.js`); базовый URL API задаётся атрибутом `data-api-base` на `` (по умолчанию `/api/v1`). | @@ -41,15 +41,16 @@ | GET | `/api/v1/stats` | Сводка для дашборда | | GET | `/api/v1/clusters` | Список кластеров | | POST | `/api/v1/clusters` | Создание в фоне (**202** + `job_id`) | -| POST | `/api/v1/clusters/{name}/start` | Запуск: **200** — `docker start` узлов (кластер в kind); **202** + `job_id` — фоновый `kind create` по сохранённому `kind-config.yaml` | -| POST | `/api/v1/clusters/{name}/stop` | Остановка узлов (`docker`/`podman` **stop**), запись в kind сохраняется | +| POST | `/api/v1/clusters/{name}/start` | Запуск в фоне (**202** + `job_id`, поле `mode`: `containers` или `kind_config`); журнал — `GET /jobs/{job_id}` | +| POST | `/api/v1/clusters/{name}/stop` | Остановка узлов в фоне (**202** + `job_id`, `mode`: `stop`); журнал — `GET /jobs/{job_id}` | | GET | `/api/v1/clusters/{name}` | Детали + `kubectl get nodes` при наличии kubeconfig | | GET | `/api/v1/clusters/{name}/kubeconfig` | Скачать файл kubeconfig | | GET | `/api/v1/clusters/{name}/workloads` | Узлы и поды (`kubectl`) | | DELETE | `/api/v1/clusters/{name}` | Удалить кластер и данные в `clusters/` | -| GET | `/api/v1/jobs` | Последние задания создания | -| GET | `/api/v1/jobs/{job_id}` | Статус одного задания (включая `progress_stage`, `progress_percent`) | -| POST | `/api/v1/jobs/{job_id}/cancel` | Запросить отмену создания (между этапами; `kind create` до конца не прерывается) | +| GET | `/api/v1/jobs` | Последние задания (`progress_log` в ответе пустой — полный журнал только в GET по `job_id`) | +| GET | `/api/v1/jobs/{job_id}` | Статус задания + полный хвост `progress_log` (лимит см. `KIND_K8S_JOB_API_LOG_MAX_LINES`) | +| DELETE | `/api/v1/jobs` | Удалить из памяти **завершённые** задания (`removed` — число записей) | +| POST | `/api/v1/jobs/{job_id}/cancel` | Прервать задание: активная команда завершается принудительно (скачивание образа, создание кластера, старт/стоп узла) | ### Фоновые задания (jobs) @@ -57,9 +58,11 @@ - В памяти держится не более **200** записей; при превышении старые задания вытесняются (`app/core/job_store.py`). - Создание кластера: `POST /api/v1/clusters` → опрос `GET /api/v1/jobs/{job_id}` (как в веб-UI). - В ответе задания поля **`progress_stage`** (текст этапа) и **`progress_percent`** (0–100) обновляются во время создания. -- Поле **`progress_log`** — массив последних строк журнала (вывод `kind create`: pull образов, подъём нод и т.д.); размер ограничен (см. `KIND_K8S_JOB_LOG_MAX_LINES` в коде `job_store`, по умолчанию до **500** строк в буфере, в JSON отдаётся хвост). -- Тип задания **`kind`**: `create_cluster` или `start_cluster` (повторный подъём по `clusters/<имя>/kind-config.yaml`). -- Статус **`cancelled`** — пользователь запросил отмену (`POST .../cancel`); этап `kind create cluster` до завершения не прерывается. +- В **GET /api/v1/jobs** (список) поле **`progress_log`** всегда **пустой массив** — меньше трафика; полный хвост — в **GET /api/v1/jobs/{job_id}** (лимит строк: `KIND_K8S_JOB_API_LOG_MAX_LINES`, по умолчанию **5000**). +- Буфер строк в памяти на задание: `KIND_K8S_JOB_LOG_MAX_LINES` (по умолчанию **2500**); при переполнении старые строки вытесняются. +- Для **`docker pull`** по умолчанию: **`--progress=plain`** и вывод **без PTY** (`KIND_K8S_DOCKER_PULL_PLAIN=1`) — в журнале идут строки с процентами/слоями. Для **podman** и **kind** — псевдо-TTY по `KIND_K8S_STREAM_PTY`, из строк убираются ANSI-коды. +- Тип задания **`kind`**: `create_cluster`, `start_cluster` (подъём по сохранённому конфигу), `start_containers` (запуск уже созданных узлов), `stop_containers` (остановка узлов). +- Статус **`cancelled`** — запрошена отмена (`POST .../cancel`); дочерний процесс текущей команды получает принудительное завершение. --- @@ -236,7 +239,7 @@ Accept: text/markdown ## GET /api/v1/jobs -Список последних фоновых заданий (создание кластера), от новых к старым. +Список последних фоновых заданий, от новых к старым. Поле **`progress_log`** в каждом элементе **пустое** — используйте **GET /api/v1/jobs/{job_id}** для журнала. **Query:** `limit` (1–200, по умолчанию **30**). @@ -253,16 +256,31 @@ Accept: text/markdown "message": "Кластер создан", "result": { "cluster_name": "dev", "kubernetes_version_tag": "v1.29.4" }, "progress_stage": null, - "progress_percent": null + "progress_percent": null, + "progress_log": [] } ] ``` --- +## DELETE /api/v1/jobs + +Удаляет из памяти записи заданий со статусом **success**, **failed** или **cancelled**. Задания **queued** и **running** не удаляются. + +**Пример ответа 200:** + +```json +{ + "removed": 4 +} +``` + +--- + ## POST /api/v1/jobs/{job_id}/cancel -Запрос отмены создания кластера. Пока задание в статусе `queued` или `running`, между этапами выполняется проверка флага; после уже запущенного `kind create cluster` нужно дождаться окончания этого шага. +Прервать фоновое задание (`create_cluster`, `start_cluster`, `start_containers`, `stop_containers`). Для длительной команды завершается связанный дочерний процесс; между шагами запуска/остановки отдельных узлов также проверяется флаг отмены. **Пример ответа 200:** @@ -270,7 +288,7 @@ Accept: text/markdown { "job_id": "a1b2…", "cancel_requested": true, - "message": "Отмена обрабатывается между этапами; во время kind create дождитесь окончания шага" + "message": "Запрошено прерывание; текущая команда будет остановлена, задание перейдёт в отменено" } ``` @@ -370,19 +388,19 @@ Accept: text/markdown ## POST /api/v1/clusters/{name}/start -Запуск кластера двумя сценариями: +Запуск кластера двумя сценариями (оба с **202** и `job_id`, журнал в `GET /jobs/{job_id}`): -1. Кластер **есть** в `kind get clusters` (узлы когда-либо создавались) — выполняется **`docker start`** / **`podman start`** для всех контейнеров с именами вида `<имя>-control-plane`, `<имя>-worker`, … Ответ **200**. -2. В **kind** кластера **нет**, но в `clusters/<имя>/kind-config.yaml` файл **есть** — ставится фоновое задание **`start_cluster`** (как при создании: `kind create` по сохранённому конфигу, журнал в `GET /jobs/{job_id}`). Ответ **202** + `job_id`. +1. Кластер **зарегистрирован** в kind — задание **`start_containers`** (поочерёдный запуск узлов, `mode`: **`containers`**). +2. В kind кластера **нет**, но есть сохранённый **`clusters/<имя>/kind-config.yaml`** — задание **`start_cluster`** (`mode`: **`kind_config`**), логика как при создании (в т.ч. скачивание образа при необходимости). -**Пример ответа 200 (контейнеры запущены):** +**Пример ответа 202 (запуск узлов):** ```json { - "name": "dev", + "job_id": "cafebabe...", + "status": "queued", "mode": "containers", - "containers_started_ok": true, - "summary": "dev-control-plane: OK; dev-worker: OK; dev-worker2: OK" + "message": "Запуск узлов; опросите GET /api/v1/jobs/{job_id}" } ``` @@ -390,9 +408,10 @@ Accept: text/markdown ```json { - "job_id": "cafebabe...", + "job_id": "deadbeef...", "status": "queued", - "message": "Подъём кластера по kind-config.yaml; опросите GET /api/v1/jobs/{job_id}" + "mode": "kind_config", + "message": "Подъём кластера по сохранённому конфигу; опросите GET /api/v1/jobs/{job_id}" } ``` @@ -402,15 +421,16 @@ Accept: text/markdown ## POST /api/v1/clusters/{name}/stop -Остановка **всех** контейнеров узлов кластера (`docker stop` / `podman stop` по префиксу имени). Запись кластера в kind **не удаляется**; позже можно снова вызвать **POST …/start** (режим `containers`). +Остановка узлов кластера **в фоне** (задание **`stop_containers`**). Запись кластера в kind **не удаляется**; позже снова **POST …/start**. -**Пример ответа 200:** +**Пример ответа 202:** ```json { - "name": "dev", - "containers_stopped_ok": true, - "summary": "dev-control-plane: OK; dev-worker: OK" + "job_id": "baba...", + "status": "queued", + "mode": "stop", + "message": "Остановка узлов; опросите GET /api/v1/jobs/{job_id}" } ``` @@ -433,13 +453,13 @@ Accept: text/markdown "created_at_utc": "2026-04-04T12:00:00+00:00", "message": null, "result": null, - "progress_stage": "kind create cluster (скачивание образов и подъём нод — может занять несколько минут)", - "progress_percent": 28, + "progress_stage": "Создание узлов кластера", + "progress_percent": 45, "progress_log": [ - "[12%] Подготовка каталога и kind-config", - "--- kind create cluster ---", - "Creating cluster \"dev\" ...", - " • Ensuring node image (kindest/node:v1.29.4) 🖼 ..." + "Подготовка конфигурации", + "Using default tag: latest", + "Status: Downloaded newer image for kindest/node:v1.29.4", + "Creating cluster \"dev\" ..." ] } ``` @@ -454,7 +474,7 @@ Accept: text/markdown "cluster_name": "dev", "created_at_utc": "2026-04-04T12:00:00+00:00", "message": "Кластер создан", - "progress_log": ["[95%] Финализация", "kubectl wait nodes: ..."], + "progress_log": ["Завершение", "Узлы готовы: condition met"], "result": { "cluster_name": "dev", "kubernetes_version_tag": "v1.29.4", diff --git a/app/models/schemas.py b/app/models/schemas.py index 2bbd71d..08f0177 100644 --- a/app/models/schemas.py +++ b/app/models/schemas.py @@ -41,11 +41,11 @@ class JobView(BaseModel): created_at_utc: str message: str | None = None result: dict[str, Any] | None = None - progress_stage: str | None = Field(default=None, description="Текущий этап создания (пока задание активно)") + progress_stage: str | None = Field(default=None, description="Текущий этап операции (пока задание активно)") progress_percent: int | None = Field(default=None, description="Прогресс 0–100 для индикатора в UI") progress_log: list[str] = Field( default_factory=list, - description="Хвост лога (kind create, этапы); обновляется при опросе GET /jobs/{id}", + description="Хвост журнала (скачивание образа, длительные команды, этапы); при опросе GET /jobs/{id}", ) diff --git a/app/static/js/dashboard.js b/app/static/js/dashboard.js index bde9807..2e64394 100644 --- a/app/static/js/dashboard.js +++ b/app/static/js/dashboard.js @@ -12,16 +12,18 @@ const body = document.body; const API = (body.dataset.apiBase || "/api/v1").replace(/\/$/, ""); - /** Интервал опроса списков и среды (мс) */ + /** Интервал автообновления таблиц и среды на панели (мс) */ var AUTO_REFRESH_MS = 3500; - /** Интервал опроса задания создания (мс) */ - var JOB_POLL_MS = 1500; + /** Опрос GET /jobs/{id} для активного задания (мс) — чаще, чтобы журнал обновлялся живее */ + var JOB_POLL_ACTIVE_MS = 650; /** @type {ReturnType | null} */ var autoTimer = null; /** @type {ReturnType | null} */ var pollTimer = null; var createInProgress = false; + /** Блокировать форму «Создать» только при создании из формы, не при старт/стоп из таблицы. */ + var pollLockCreateForm = false; /** @type {string | null} */ var currentPollJobId = null; /** Имя кластера в открытой модалке «Состояние» (для скачивания kubeconfig). */ @@ -69,6 +71,32 @@ return d.innerHTML; } + function pad2(n) { + return n < 10 ? "0" + n : String(n); + } + + /** + * Краткий человекочитаемый вид даты/времени в UTC для таблицы заданий (ДД.ММ.ГГ ЧЧ:ММ UTC). + * @param {string} iso + */ + function formatJobTimeUtc(iso) { + if (!iso) return "—"; + const d = new Date(iso); + if (Number.isNaN(d.getTime())) return iso; + return ( + pad2(d.getUTCDate()) + + "." + + pad2(d.getUTCMonth() + 1) + + "." + + String(d.getUTCFullYear()).slice(-2) + + " " + + pad2(d.getUTCHours()) + + ":" + + pad2(d.getUTCMinutes()) + + " UTC" + ); + } + /** * Скачать kubeconfig кластера (GET /clusters/{name}/kubeconfig). * @param {string} clusterName @@ -493,6 +521,43 @@ } } + /** + * Удалить из памяти сервера только завершённые задания (успех / ошибка / отмена). + */ + async function clearCompletedJobs() { + const ok = await openConfirmModal({ + title: "Очистить завершённые задания?", + message: + "Из списка на сервере будут удалены только завершённые записи. Текущее активное задание не затрагивается.", + confirmLabel: "Очистить", + danger: false, + }); + if (!ok) return; + const msg = document.getElementById("jobs-msg"); + try { + const r = await fetch(API + "/jobs", { method: "DELETE", headers: { Accept: "application/json" } }); + const text = await r.text(); + var data = {}; + try { + data = text ? JSON.parse(text) : {}; + } catch (e) { + data = {}; + } + if (!r.ok) { + showToast(formatApiError(data, text || r.statusText), true); + return; + } + var n = data.removed != null ? data.removed : 0; + if (msg) msg.textContent = n ? "Удалено записей: " + n : "Нечего очищать."; + showToast(n ? "Удалено завершённых заданий: " + n : "Нет завершённых заданий для очистки", false); + await loadJobs(); + await loadStats(); + } catch (e) { + if (msg) msg.textContent = "Ошибка: " + e.message; + showToast(e.message, true); + } + } + async function loadJobs() { const tbody = document.querySelector("#tbl-jobs tbody"); const msg = document.getElementById("jobs-msg"); @@ -503,7 +568,10 @@ rows.forEach(function (j) { const tr = document.createElement("tr"); const st = escapeHtml(j.status || ""); - var kindTag = j.kind === "start_cluster" ? "[старт] " : ""; + var kindTag = ""; + if (j.kind === "start_cluster") kindTag = "[старт] "; + else if (j.kind === "stop_containers") kindTag = "[стоп] "; + else if (j.kind === "start_containers") kindTag = "[запуск узлов] "; var cellMsg = kindTag + (j.message || "").slice(0, 140); if ((j.status === "running" || j.status === "queued") && j.progress_stage) { cellMsg = @@ -511,11 +579,15 @@ j.progress_stage + (j.progress_percent != null ? " (" + j.progress_percent + "%)" : ""); } + const iso = j.created_at_utc || ""; + const timeLabel = formatJobTimeUtc(iso); tr.innerHTML = "" + "" + escapeHtml(j.cluster_name || "—") + @@ -644,14 +716,31 @@ message: "Кластер «" + name + - "»: контейнеры будут остановлены (docker/podman stop). Запись в kind сохранится — позже можно снова нажать «Старт».", + "»: узлы будут остановлены. Сведения о кластере сохранятся — позже можно снова нажать «Старт».", confirmLabel: "Остановить", danger: false, }); if (!ok) return; + const url = API + "/clusters/" + encodeURIComponent(name) + "/stop"; try { - const res = await api("/clusters/" + encodeURIComponent(name) + "/stop", { method: "POST" }); - showToast(String(res.summary || "Узлы остановлены"), false); + const r = await fetch(url, { method: "POST", headers: { Accept: "application/json" } }); + const text = await r.text(); + var data = {}; + try { + data = text ? JSON.parse(text) : {}; + } catch (e) { + data = {}; + } + if (r.status === 202 && data.job_id) { + setProgressHint(name, "stop"); + pollJob(data.job_id, { lockCreateForm: false }); + return; + } + if (!r.ok) { + showToast(formatApiError(data, text || r.statusText), true); + return; + } + showToast(String(data.summary || "Готово"), false); await loadClusters(); await loadStats(); await loadHealth(); @@ -674,8 +763,10 @@ data = {}; } if (r.status === 202 && data.job_id) { - setProgressHint(name); - pollJob(data.job_id); + var sm = data.mode || "kind_config"; + if (sm === "containers") setProgressHint(name, "start_containers"); + else setProgressHint(name, "start_config"); + pollJob(data.job_id, { lockCreateForm: false }); return; } if (!r.ok) { @@ -741,21 +832,26 @@ } /** - * Обновить текст подсказки над прогрессом (создание с нуля или старт по конфигу). + * Подсказка над прогресс-баром (без технических имён команд). * @param {string | null} clusterName + * @param {"create"|"start_config"|"start_containers"|"stop"} [mode] */ - function setProgressHint(clusterName) { + function setProgressHint(clusterName, mode) { const hint = document.getElementById("create-progress-hint"); if (!hint) return; - if (clusterName) { - hint.innerHTML = - "Кластер " + - escapeHtml(clusterName) + - ": подъём по конфигу или длительный kind create — ниже журнал (pull образов и ноды)."; + mode = mode || "create"; + var named = clusterName ? "«" + clusterName + "»" : ""; + if (mode === "stop") { + hint.textContent = "Останавливаем кластер " + named + ". Ход операции — в журнале ниже."; + } else if (mode === "start_containers") { + hint.textContent = "Запускаем узлы кластера " + named + ". Ход операции — в журнале ниже."; + } else if (mode === "start_config") { + hint.textContent = "Поднимаем кластер " + named + " по сохранённой конфигурации. Журнал ниже."; } else { - hint.innerHTML = - "Создание кластера: шаг kind create может занять несколько минут при первом pull образов. " + - "Ниже — журнал в реальном времени."; + hint.textContent = + "Создаём кластер" + + (named ? " " + named : "") + + ". Сначала при необходимости скачивается образ, затем создаются узлы — всё в журнале ниже."; } } @@ -781,7 +877,10 @@ } currentPollJobId = null; createInProgress = false; - setCreateFormDisabled(false); + if (pollLockCreateForm) { + setCreateFormDisabled(false); + pollLockCreateForm = false; + } showCreateProgress(false); } @@ -789,13 +888,19 @@ if (!currentPollJobId) return; try { await api("/jobs/" + encodeURIComponent(currentPollJobId) + "/cancel", { method: "POST" }); - showToast("Запрос отмены отправлен (между этапами)", false); + showToast("Запрос на прерывание отправлен", false); } catch (e) { showToast(e.message || "Не удалось отменить", true); } } - function pollJob(jobId) { + /** + * @param {string} jobId + * @param {{ lockCreateForm?: boolean }} [opts] + */ + function pollJob(jobId, opts) { + opts = opts || {}; + pollLockCreateForm = !!opts.lockCreateForm; const details = document.getElementById("job-details"); const msg = document.getElementById("create-msg"); const logEl = document.getElementById("job-log-panel"); @@ -807,7 +912,9 @@ if (pollTimer) clearInterval(pollTimer); createInProgress = true; currentPollJobId = jobId; - setCreateFormDisabled(true); + if (pollLockCreateForm) { + setCreateFormDisabled(true); + } showCreateProgress(true); updateCreateProgressFromJob({ progress_percent: 0, progress_stage: "В очереди…", status: "queued" }); @@ -817,23 +924,35 @@ const preEl = document.getElementById("job-json"); if (preEl) preEl.textContent = JSON.stringify(j, null, 2); updateCreateProgressFromJob(j); - if (logEl && j.progress_log && j.progress_log.length) { - logEl.textContent = j.progress_log.join("\n"); - logEl.scrollTop = logEl.scrollHeight; + if (logEl) { + if (j.progress_log && j.progress_log.length) { + logEl.removeAttribute("data-placeholder"); + logEl.textContent = j.progress_log.join("\n"); + logEl.scrollTop = logEl.scrollHeight; + } else if (j.status === "running" || j.status === "queued") { + logEl.setAttribute("data-placeholder", "1"); + logEl.textContent = + "Журнал подгружается… Строки скачивания и создания узлов появляются по мере вывода команды."; + } } if (j.status === "success" || j.status === "failed" || j.status === "cancelled") { stopPollJob(); - setProgressHint(null); + setProgressHint(null, "create"); if (msg) { if (j.status === "success") { - msg.textContent = - j.kind === "start_cluster" ? "Кластер поднят по сохранённому конфигу." : "Кластер создан."; + if (j.kind === "start_cluster") msg.textContent = "Кластер поднят по сохранённому конфигу."; + else if (j.kind === "stop_containers") msg.textContent = "Узлы кластера остановлены."; + else if (j.kind === "start_containers") msg.textContent = "Узлы кластера запущены."; + else msg.textContent = "Кластер создан."; } else if (j.status === "cancelled") msg.textContent = j.message || "Операция отменена."; else msg.textContent = "Ошибка: " + (j.message || ""); } if (j.status === "success") { - showToast(j.kind === "start_cluster" ? "Кластер запущен" : "Кластер создан", false); + if (j.kind === "start_cluster") showToast("Кластер запущен", false); + else if (j.kind === "stop_containers") showToast("Узлы остановлены", false); + else if (j.kind === "start_containers") showToast("Узлы запущены", false); + else showToast("Кластер создан", false); } else if (j.status === "cancelled") showToast(j.message || "Отменено", false); else showToast(j.message || "Ошибка", true); await loadClusters(); @@ -846,7 +965,7 @@ } }; tick(); - pollTimer = setInterval(tick, JOB_POLL_MS); + pollTimer = setInterval(tick, JOB_POLL_ACTIVE_MS); } function refreshLists() { @@ -866,7 +985,7 @@ const details = document.getElementById("job-details"); if (msg) msg.textContent = ""; if (details) details.classList.add("hidden"); - setProgressHint(null); + setProgressHint(null, "create"); const fd = new FormData(form); const body = { name: String(fd.get("name") || "").trim(), @@ -880,7 +999,7 @@ body: JSON.stringify(body), }); if (msg) msg.textContent = "Задание: " + res.job_id; - pollJob(res.job_id); + pollJob(res.job_id, { lockCreateForm: true }); } catch (e) { if (msg) msg.textContent = "Ошибка: " + e.message; showToast(e.message, true); @@ -895,6 +1014,13 @@ }); } + const btnClearJobs = document.getElementById("btn-jobs-clear"); + if (btnClearJobs) { + btnClearJobs.addEventListener("click", function () { + clearCompletedJobs(); + }); + } + const mClose = document.getElementById("modal-close"); if (mClose) mClose.addEventListener("click", closeModal); diff --git a/app/static/style.css b/app/static/style.css index c177a02..202eacb 100644 --- a/app/static/style.css +++ b/app/static/style.css @@ -418,7 +418,7 @@ body.modal-open { } .job-log-panel { margin: 0; - max-height: min(40vh, 280px); + max-height: min(55vh, 420px); overflow: auto; padding: 0.5rem 0.65rem; font-size: 0.78rem; @@ -430,6 +430,16 @@ body.modal-open { word-break: break-word; } +/* Заголовок блока заданий + кнопка очистки */ +.jobs-card-toolbar { + margin-bottom: 0.65rem; +} +.jobs-card-title { + margin: 0; + font-size: 1.15rem; + line-height: 1.3; +} + .git-hint { font-size: 0.85rem; margin: 0 0 0.65rem; @@ -690,6 +700,12 @@ table { font-size: 0.9rem; } +th .th-hint { + font-weight: 500; + opacity: 0.7; + font-size: 0.82em; +} + th, td { border-bottom: 1px solid var(--border); diff --git a/app/templates/dashboard.html b/app/templates/dashboard.html index 0ba137c..3963112 100644 --- a/app/templates/dashboard.html +++ b/app/templates/dashboard.html @@ -111,7 +111,7 @@ aria-relevant="additions" > - +

@@ -146,12 +146,15 @@
-

Последние задания

+
+

Последние задания

+ +
- + diff --git a/docker-compose.yml b/docker-compose.yml index 34b9f86..646538c 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -40,6 +40,11 @@ services: KIND_K8S_HUB_TAGS_MAX_PAGES: ${KIND_K8S_HUB_TAGS_MAX_PAGES:-} KIND_K8S_DEBUG: ${KIND_K8S_DEBUG:-} KIND_K8S_JOB_LOG_MAX_LINES: ${KIND_K8S_JOB_LOG_MAX_LINES:-} + # Псевдо-TTY для потоковых команд (0 = pipe). Для docker pull по умолчанию используется --progress=plain без PTY. + KIND_K8S_STREAM_PTY: ${KIND_K8S_STREAM_PTY:-} + KIND_K8S_DOCKER_PULL_PLAIN: ${KIND_K8S_DOCKER_PULL_PLAIN:-} + # Сколько строк журнала отдавать в GET /api/v1/jobs/{id} (по умолчанию 5000 в коде). + KIND_K8S_JOB_API_LOG_MAX_LINES: ${KIND_K8S_JOB_API_LOG_MAX_LINES:-} KIND_K8S_README_PATH: ${KIND_K8S_README_PATH:-} KIND_K8S_WAIT_NODES: ${KIND_K8S_WAIT_NODES:-} KIND_K8S_WAIT_NODES_TIMEOUT_SEC: ${KIND_K8S_WAIT_NODES_TIMEOUT_SEC:-}
Время (UTC)Дата, время (UTC) Кластер Статус Сообщение