From 7cadd1d6735a98488b501af01eab3f7d191b5d2a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=94=D0=B0=D0=BD=D0=B8=D0=B8=D0=BB=20=D0=93=D1=80=D0=B0?= =?UTF-8?q?=D0=B1=D0=B0=D1=80=D1=8C?= Date: Fri, 22 May 2026 18:30:34 +0300 Subject: [PATCH] fix: enhance subprocess management and progress reporting in setup and update commands --- cobot/commands/local_setup.py | 40 ++++++++++++++++-------- cobot/commands/update.py | 30 +++++++++++------- cobot/tui.py | 59 ++++++++++++++++++++++++++++------- 3 files changed, 92 insertions(+), 37 deletions(-) diff --git a/cobot/commands/local_setup.py b/cobot/commands/local_setup.py index 9182933..5a20a84 100644 --- a/cobot/commands/local_setup.py +++ b/cobot/commands/local_setup.py @@ -250,7 +250,11 @@ def _run_apt_with_progress( # Installation steps - each step maps to one visible phase in the log screen # Шаги установки - каждый шаг соответствует одной видимой фазе в экране лога -def _step_prereqs(write: Write, on_progress: Callable[[float], None]) -> None: +def _step_prereqs( + write: Write, + on_progress: Callable[[float], None], + register_proc: Callable | None = None, +) -> None: """Step 1 - Refresh apt cache and install packages required for the ROS2 setup. Installs: software-properties-common, curl, gnupg2, lsb-release, build-essential. @@ -264,7 +268,7 @@ def _step_prereqs(write: Write, on_progress: Callable[[float], None]) -> None: write("[cyan][*][/cyan] Updating package lists...") _run_apt_with_progress( ["sudo", "apt-get", "update"] + _APT_OPTS, - write, on_progress, _APT_ENV, + write, on_progress, _APT_ENV, register_proc, ) write("[cyan][*][/cyan] Installing prerequisites...") _run_apt_with_progress( @@ -273,18 +277,23 @@ def _step_prereqs(write: Write, on_progress: Callable[[float], None]) -> None: "software-properties-common", "curl", "gnupg2", "lsb-release", "build-essential", ] + _APT_OPTS, - write, on_progress, _APT_ENV, + write, on_progress, _APT_ENV, register_proc, ) write("[cyan][*][/cyan] Adding universe repository...") # --no-update prevents add-apt-repository from running its own apt-get update, # which would ignore our timeout options and could hang indefinitely. # --no-update запрещает add-apt-repository запускать собственный apt-get update, # который игнорирует наши таймауты и может зависнуть. - _run_logged(["sudo", "add-apt-repository", "-y", "--no-update", "universe"], write) + _run_logged(["sudo", "add-apt-repository", "-y", "--no-update", "universe"], write, + register_proc=register_proc) write("[green][ok][/green] Prerequisites ready") -def _step_ros2_repo(write: Write, on_progress: Callable[[float], None]) -> None: +def _step_ros2_repo( + write: Write, + on_progress: Callable[[float], None], + register_proc: Callable | None = None, +) -> None: """Step 2 - Download the ROS2 signing key and register the ROS2 apt repository. Removes any previous key and sources file first so re-runs always start clean. @@ -351,7 +360,7 @@ def _step_ros2_repo(write: Write, on_progress: Callable[[float], None]) -> None: "-o", "Dir::Etc::sourceparts=-", "-o", "APT::Get::List-Cleanup=0", ] + _APT_OPTS, - write, on_progress, _APT_ENV, + write, on_progress, _APT_ENV, register_proc, ) write("[green][ok][/green] ROS2 repository ready") @@ -380,7 +389,11 @@ def _step_install_ros2( write(f"[green][ok][/green] ros-{_DISTRO}-{pkg} installed") -def _step_dev_tools(write: Write, on_progress: Callable[[float], None]) -> None: +def _step_dev_tools( + write: Write, + on_progress: Callable[[float], None], + register_proc: Callable | None = None, +) -> None: """Step 4 - Install colcon, rosdep, vcstool and initialize rosdep. Installs the Python packages needed to build and manage ROS2 workspaces. @@ -403,15 +416,16 @@ def _step_dev_tools(write: Write, on_progress: Callable[[float], None]) -> None: "python3-rosdep", "python3-vcstool", ] + _APT_OPTS, - write, on_progress, _APT_ENV, + write, on_progress, _APT_ENV, register_proc, ) if not _ROSDEP_SOURCES.exists(): write("[cyan][*][/cyan] Initializing rosdep...") - _run_logged(["sudo", "rosdep", "init"], write) + _run_logged(["sudo", "rosdep", "init"], write, register_proc=register_proc) else: write("[green][ok][/green] rosdep already initialized") write("[cyan][*][/cyan] Updating rosdep...") - _run_logged(["rosdep", "update", "--rosdistro", _DISTRO], write) + _run_logged(["rosdep", "update", "--rosdistro", _DISTRO], write, + register_proc=register_proc) write("[green][ok][/green] Dev tools ready") @@ -474,12 +488,12 @@ def _task_install(screen: LogScreen, pkg: str) -> None: return lambda p: screen.set_progress(lo + p / 100.0 * (hi - lo)) screen.set_progress(0, "Preparing...") - _step_prereqs(screen.write, prog(0, 15)) + _step_prereqs(screen.write, prog(0, 15), register_proc=screen.set_proc) if screen.is_stopped(): return screen.set_progress(15, "Setting up ROS2 repository...") - _step_ros2_repo(screen.write, prog(15, 30)) + _step_ros2_repo(screen.write, prog(15, 30), register_proc=screen.set_proc) if screen.is_stopped(): return @@ -489,7 +503,7 @@ def _task_install(screen: LogScreen, pkg: str) -> None: return screen.set_progress(75, "Installing dev tools...") - _step_dev_tools(screen.write, prog(75, 95)) + _step_dev_tools(screen.write, prog(75, 95), register_proc=screen.set_proc) if screen.is_stopped(): return diff --git a/cobot/commands/update.py b/cobot/commands/update.py index 3cc7768..71608d3 100644 --- a/cobot/commands/update.py +++ b/cobot/commands/update.py @@ -33,14 +33,16 @@ def _task_update(screen: LogScreen) -> None: # Fetch (0 → 30 %) screen.set_progress(0, "Fetching from remote...") screen.write("[cyan][*][/cyan] Fetching from remote...") - fetch = subprocess.run( + fetch_proc = subprocess.Popen( ["git", "fetch", "origin"], - cwd=_PROJECT_DIR, capture_output=True, text=True, + cwd=_PROJECT_DIR, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, ) + screen.set_proc(fetch_proc) + fetch_out, fetch_err = fetch_proc.communicate() if screen.is_stopped(): return - if fetch.returncode != 0: - screen.write(f"[red]Fetch failed:[/red] {fetch.stderr.strip()}") + if fetch_proc.returncode not in (0, -9): + screen.write(f"[red]Fetch failed:[/red] {fetch_err.strip()}") screen.finish(False) return screen.set_progress(30) @@ -72,14 +74,16 @@ def _task_update(screen: LogScreen) -> None: # Pull (30 → 80 %) screen.set_progress(30, "Pulling changes...") screen.write("\n[cyan][*][/cyan] Pulling changes...") - pull = subprocess.run( + pull_proc = subprocess.Popen( ["git", "pull", "origin", branch], - capture_output=True, text=True, cwd=_PROJECT_DIR, + stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, cwd=_PROJECT_DIR, ) + screen.set_proc(pull_proc) + pull_out, pull_err = pull_proc.communicate() if screen.is_stopped(): return - if pull.returncode != 0: - for line in (pull.stdout + pull.stderr).splitlines(): + if pull_proc.returncode not in (0, -9): + for line in (pull_out + pull_err).splitlines(): if line.strip(): screen.write(line) screen.write("[red]Pull failed.[/red]") @@ -92,16 +96,18 @@ def _task_update(screen: LogScreen) -> None: # Переустанавливаем, чтобы бинарник cobot подхватил новые зависимости из pyproject.toml. screen.set_progress(80, "Reinstalling cobot CLI...") screen.write("\n[cyan][*][/cyan] Reinstalling cobot CLI...") - reinstall = subprocess.run( + reinstall_proc = subprocess.Popen( ["uv", "tool", "install", "--editable", str(_PROJECT_DIR)], - capture_output=True, text=True, + stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, ) + screen.set_proc(reinstall_proc) + reinstall_out, reinstall_err = reinstall_proc.communicate() if screen.is_stopped(): return - if reinstall.returncode == 0: + if reinstall_proc.returncode in (0, -9): screen.write("[green][ok][/green] cobot reinstalled") else: - screen.write(f"[yellow]Warning:[/yellow] reinstall failed — {reinstall.stderr.strip()}") + screen.write(f"[yellow]Warning:[/yellow] reinstall failed — {reinstall_err.strip()}") if not screen.is_stopped(): screen.set_progress(100, "Done") diff --git a/cobot/tui.py b/cobot/tui.py index d220ea6..706cfb9 100644 --- a/cobot/tui.py +++ b/cobot/tui.py @@ -1,5 +1,7 @@ from __future__ import annotations +import os +import signal from typing import Callable, List, Optional from textual import on @@ -195,9 +197,10 @@ class LogScreen(Screen[bool]): self._finished = False self._success = False self._show_progress = show_progress - # Tracks the subprocess that is currently running so on_unmount can kill it. - # Отслеживает текущий subprocess, чтобы on_unmount мог его завершить. + # All subprocesses registered via set_proc() - every one gets killed on unmount. + # Все подпроцессы зарегистрированные через set_proc() - каждый убивается при выходе. self._active_proc = None + self._procs: list = [] self._stopped = False def compose(self) -> ComposeResult: @@ -217,11 +220,12 @@ class LogScreen(Screen[bool]): self.app.run_worker(lambda: self._run_fn(self), thread=True) def set_proc(self, proc) -> None: - # Register the subprocess that is currently running. - # Called from the worker thread - GIL makes simple assignment safe here. - # Регистрируем текущий subprocess. - # Вызывается из рабочего потока - простое присваивание безопасно благодаря GIL. + # Register the subprocess that is currently running. Added to _procs so on_unmount + # can kill it even if another proc is registered afterwards. + # Регистрируем текущий subprocess. Добавляем в _procs, чтобы on_unmount мог его убить + # даже если после него будет зарегистрирован другой процесс. self._active_proc = proc + self._procs.append(proc) def is_stopped(self) -> bool: # Return True if the user has closed the screen before the task finished. @@ -229,15 +233,24 @@ class LogScreen(Screen[bool]): return self._stopped def on_unmount(self) -> None: - # Kill the active subprocess when the screen closes so it does not keep running in the background. - # Убиваем активный subprocess при закрытии экрана, чтобы он не продолжал работать в фоне. + # Kill every registered subprocess so nothing keeps running in the background after exit. + # Use SIGKILL on the process group to also terminate any children spawned by the process + # (e.g. dpkg or apt subprocesses spawned under sudo). Falls back to proc.kill() if the + # process group is not available (e.g. already exited). + # Убиваем все зарегистрированные подпроцессы, чтобы ничего не висело в фоне после выхода. + # Используем SIGKILL по группе процессов, чтобы завершить и дочерние процессы + # (например dpkg или apt запущенные под sudo). Откат на proc.kill() если группа недоступна. self._stopped = True - proc = self._active_proc - if proc is not None: + for proc in list(self._procs): try: - proc.kill() + pgid = os.getpgid(proc.pid) + os.killpg(pgid, signal.SIGKILL) except Exception: - pass + try: + proc.kill() + except Exception: + pass + self._procs.clear() def set_progress(self, pct: float, label: str = "") -> None: # Thread-safe - this is called from the worker thread, not the UI thread. @@ -305,6 +318,7 @@ class RunScreen(Screen[None]): self._kill_fn = None # optional custom kill callable, set via set_kill_fn() self._finished = False self._stopped = False + self._procs: list = [] # all registered procs for cleanup on forced exit def compose(self) -> ComposeResult: yield Static(self._title, id="step") @@ -323,6 +337,7 @@ class RunScreen(Screen[None]): # Register the subprocess so the Stop button knows what to terminate. # Регистрируем subprocess, чтобы кнопка Stop знала что завершать. self._proc = proc + self._procs.append(proc) def set_kill_fn(self, fn: Callable) -> None: # Override the default proc.terminate() with a custom kill function. @@ -353,6 +368,26 @@ class RunScreen(Screen[None]): msg = "[green]Process exited.[/green] Press [bold]Enter[/bold] to close." self.query_one("#hint", Static).update(msg) + def on_unmount(self) -> None: + # Kill all registered subprocesses when the screen is forcibly closed (e.g. Ctrl+Q). + # Убиваем все зарегистрированные подпроцессы при принудительном закрытии экрана (Ctrl+Q). + self._stopped = True + if self._kill_fn is not None: + try: + self._kill_fn() + except Exception: + pass + for proc in list(self._procs): + try: + pgid = os.getpgid(proc.pid) + os.killpg(pgid, signal.SIGKILL) + except Exception: + try: + proc.kill() + except Exception: + pass + self._procs.clear() + def action_stop_close(self) -> None: # This runs in the UI thread, so we call _append() directly instead of write() # because write() uses call_from_thread() which only works from other threads.