From 426bac85bf5e1d4d5946107f8128aa59629897bf Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 21:16:03 -0700 Subject: [PATCH] refactor(hermes_cli/update): split Windows pause/resume/venv-guard god functions into phase helpers; unify psutil/sc.exe/reap-rescan patterns; compact restart_recovery bucketing --- hermes_cli/update_cmd_windows.py | 692 +++++++++++++------------- hermes_cli/update_restart_recovery.py | 67 +-- 2 files changed, 367 insertions(+), 392 deletions(-) diff --git a/hermes_cli/update_cmd_windows.py b/hermes_cli/update_cmd_windows.py index 2933ccbdf4..307d7a629f 100644 --- a/hermes_cli/update_cmd_windows.py +++ b/hermes_cli/update_cmd_windows.py @@ -53,24 +53,19 @@ def _wait_for_windows_update_gateway_exit( from gateway.status import _pid_exists + def _alive(pid: int) -> bool: + try: + return bool(_pid_exists(pid)) + except Exception: + return False + remaining = set(pids) deadline = _time.monotonic() + max(timeout, 0.0) while remaining and _time.monotonic() < deadline: - for pid in list(remaining): - try: - if not _pid_exists(pid): - remaining.discard(pid) - except Exception: - remaining.discard(pid) + remaining = {pid for pid in remaining if _alive(pid)} if remaining: _time.sleep(0.25) - - survivors: set[int] = set() - for pid in remaining: - with suppress(Exception): - if _pid_exists(pid): - survivors.add(pid) - return survivors + return {pid for pid in remaining if _alive(pid)} def _self_and_non_gateway_ancestor_pids(psutil) -> set[int]: @@ -87,16 +82,21 @@ def _self_and_non_gateway_ancestor_pids(psutil) -> set[int]: skip: set[int] = {os.getpid()} with suppress(Exception): for anc in psutil.Process().parents(): - try: - anc_cmdline = " ".join(anc.cmdline() or []) - except Exception: - anc_cmdline = "" + anc_cmdline = _cmdline_or_empty(anc) if _is_gw is not None and anc_cmdline and _is_gw(anc_cmdline): continue skip.add(int(anc.pid)) return skip +def _cmdline_or_empty(proc) -> str: + """Joined argv of a psutil process, ``""`` when it can't be read.""" + try: + return " ".join(proc.cmdline() or []) + except Exception: + return "" + + def _lower_dir_prefix(path: Path) -> str: """``str(path)`` lower-cased with one trailing separator, resolved when possible (prefix matching).""" try: @@ -106,6 +106,21 @@ def _lower_dir_prefix(path: Path) -> str: return raw.lower().rstrip(os.sep) + os.sep +def _psutil(): + """The ``psutil`` module, or ``None`` when it can't be imported (callers degrade, never raise).""" + try: + import psutil # type: ignore + except Exception: + return None + return psutil + + +def _parent_is_live(proc) -> bool: + """True when *proc* has a running parent that is not a recycled PID (a "parent" created after its child).""" + parent = proc.parent() + return parent is not None and parent.is_running() and parent.create_time() <= proc.create_time() + + def _detect_venv_python_processes( *, exclude_pids: set[int] | None = None ) -> list[tuple[int, str, str]]: @@ -116,18 +131,13 @@ def _detect_venv_python_processes( respawns its backend) so callers should refuse. Empty off-Windows / without psutil; self+ancestors excluded. """ from hermes_cli.update_cmd import _m - if not _m()._is_windows(): - return [] - try: - import psutil - except Exception: + psutil = _psutil() + if not _m()._is_windows() or psutil is None: return [] venv_prefix = _lower_dir_prefix(_m().PROJECT_ROOT / "venv") root_prefix = _lower_dir_prefix(_m().PROJECT_ROOT) - - skip: set[int] = set(exclude_pids or set()) - skip |= _self_and_non_gateway_ancestor_pids(psutil) + skip = set(exclude_pids or set()) | _self_and_non_gateway_ancestor_pids(psutil) matches: list[tuple[int, str, str]] = [] try: @@ -160,10 +170,7 @@ def _detect_venv_python_processes( ): continue - try: - cmdline_raw = " ".join(proc.cmdline() or []) - except Exception: - cmdline_raw = "" + cmdline_raw = _cmdline_or_empty(proc) cmdline_low = cmdline_raw.lower() # Fallback: uv/base-interpreter trampolines have an exe OUTSIDE the venv yet hold # its .pyd files — match cmdline (venv path, or `-m hermes_cli.main` + root/cwd). @@ -178,10 +185,9 @@ def _detect_venv_python_processes( is_holder = True if not is_holder: continue - name = info.get("name") or Path(exe).name # FULL cmdline: callers parse it (pausable-gateway exemption looks for `gateway run`); # truncating here misreported autostarted gateways as blockers. Truncate at display time. - matches.append((int(pid), str(name), cmdline_raw)) + matches.append((int(pid), name, cmdline_raw)) return matches @@ -299,17 +305,12 @@ def _venv_launcher_ancestors(pids: list[int]) -> list[int]: paused gateway still tripped the guard. One hop up only, venv-prefixed only (bounds blast radius). """ from hermes_cli.update_cmd import _m - if not _m()._is_windows() or not pids: - return [] - try: - import psutil - except Exception: + psutil = _psutil() + if not _m()._is_windows() or not pids or psutil is None: return [] venv_prefix = _lower_dir_prefix(_m().PROJECT_ROOT / "venv") - skip = _self_and_non_gateway_ancestor_pids(psutil) - found: list[int] = [] for pid in pids: try: @@ -341,11 +342,7 @@ def _leftover_pausable_gateway_pids( """ from hermes_cli._scan_venv_blockers import _is_pausable_gateway - try: - import psutil # type: ignore - except Exception: - psutil = None - + psutil = _psutil() pids: list[int] = [] for pid, _name, cmdline in matches: argv = cmdline @@ -406,17 +403,14 @@ def _ledger_manual_serve_holders( except Exception: return [] holder_pids = {int(pid) for pid, _name, _cmd in matches} - out: list[dict] = [] - for entry in ledger_entries(): - if entry.get("purpose") not in ("serve", "dashboard"): - continue - pid = entry.get("pid") - if not isinstance(pid, int) or pid not in holder_pids: - continue - if spawner_is_dead(entry) is False: - continue # live Desktop supervisor owns it — keep refusing - out.append(entry) - return out + return [ + entry + for entry in ledger_entries() + if entry.get("purpose") in ("serve", "dashboard") + and isinstance(entry.get("pid"), int) + and entry["pid"] in holder_pids + and spawner_is_dead(entry) is not False # False = live Desktop supervisor owns it; keep refusing + ] def _serve_relaunch_commands(entries: list[dict]) -> list[list[str]]: @@ -425,19 +419,14 @@ def _serve_relaunch_commands(entries: list[dict]) -> list[list[str]]: """ from hermes_cli.update_cmd import _m commands: list[list[str]] = [] - hermes = None - try: + hermes = "hermes" + with suppress(Exception): scripts_dir = _m()._venv_scripts_dir() if scripts_dir is not None: for name in ("hermes.exe", "hermes"): - candidate = scripts_dir / name - if candidate.is_file(): - hermes = str(candidate) + if (scripts_dir / name).is_file(): + hermes = str(scripts_dir / name) break - except Exception: - hermes = None - if hermes is None: - hermes = "hermes" for entry in entries: port = entry.get("port") if not isinstance(port, int) or port <= 0: @@ -515,9 +504,8 @@ def _orphaned_desktop_backend_pids( accepted root's tree fold into it; only roots are returned (``taskkill /T`` reaps descendants). Any other live-parent backend, unjustified non-backend, unprovable case, or no psutil -> ``None``. Never raises. """ - try: - import psutil # type: ignore - except Exception: + psutil = _psutil() + if psutil is None: return None # Pass 1: find orphaned backend ROOTS among the holders. @@ -543,13 +531,12 @@ def _orphaned_desktop_backend_pids( try: ppid = proc.ppid() parent = psutil.Process(ppid) if ppid else None - if parent is not None and parent.is_running(): - # PID-reuse check: a "parent" created after its child is a recycled PID. - if parent.create_time() <= proc.create_time(): - # Live parent: not a root, maybe an orphan root's descendant (the venv - # trampoline re-execs uv python with the SAME argv). Defer to pass 2. - remaining.append((int(pid), low)) - continue + # PID-reuse check: a "parent" created after its child is a recycled PID. + if parent is not None and parent.is_running() and parent.create_time() <= proc.create_time(): + # Live parent: not a root, maybe an orphan root's descendant (the venv + # trampoline re-execs uv python with the SAME argv). Defer to pass 2. + remaining.append((int(pid), low)) + continue except psutil.NoSuchProcess: pass # parent gone → orphan except Exception: @@ -593,16 +580,13 @@ def _ledger_reapable_backend_pids( except Exception: return [] by_pid = {e.get("pid"): e for e in entries if isinstance(e.get("pid"), int)} - roots: list[int] = [] - for pid, _name, _cmdline in matches: - entry = by_pid.get(int(pid)) - if not entry: - continue - if entry.get("purpose") not in REAPABLE_PURPOSES: - continue - if spawner_is_dead(entry) is True: - roots.append(int(pid)) - return roots + return [ + int(pid) + for pid, _name, _cmdline in matches + if (entry := by_pid.get(int(pid))) + and entry.get("purpose") in REAPABLE_PURPOSES + and spawner_is_dead(entry) is True + ] def _handoff_reapable_backend_pids( @@ -616,11 +600,9 @@ def _handoff_reapable_backend_pids( leaks. Only Hermes backends qualify — any non-backend holder, or no psutil -> ``None``. The CALLER must have confirmed the hand-off gate; outside it the stricter orphan-only path stands. """ - try: - import psutil # type: ignore - except Exception: + psutil = _psutil() + if psutil is None: return None - roots: list[int] = [] for pid, _name, cmdline in matches: low = _live_argv_low(psutil, pid, cmdline) @@ -629,7 +611,6 @@ def _handoff_reapable_backend_pids( if not _is_backend_argv(low): return None # unexpected non-backend holder: refuse the whole set roots.append(int(pid)) - return roots or None @@ -699,11 +680,7 @@ def _desktop_owns_gateway_lifecycle() -> bool: if spawner_is_dead(entry) is False: return True - try: - import psutil - except Exception: - psutil = None - + psutil = _psutil() try: holders = _m()._detect_venv_python_processes() except Exception as exc: @@ -715,31 +692,39 @@ def _desktop_owns_gateway_lifecycle() -> bool: continue if psutil is None: return True # cannot prove orphanhood; a live control plane suffices - try: - proc = psutil.Process(int(pid)) - parent = proc.parent() - if parent is None or not parent.is_running(): - continue - if parent.create_time() > proc.create_time(): - continue - return True - except Exception: - continue + with suppress(Exception): + if _parent_is_live(psutil.Process(int(pid))): + return True return False -def _stop_windows_gateway_service( - name: str, - *, - expected_processes: tuple[tuple[int, float], ...] = (), - expected_service_identity: tuple[int, float] | None = None, - expected_gateway_identity: tuple[int, float] | None = None, - timeout: float = 30.0, -) -> None: - """Stop one verified Windows service and wait until SCM reports it down.""" - import psutil # noqa: PLC0415 +def _sc_exe(verb: str, name: str, service, settled_status: str) -> None: + """``sc.exe ``; a non-zero exit is only an error when SCM doesn't already report *settled_status*.""" + result = subprocess.run( + ["sc.exe", verb, name], + capture_output=True, + text=True, + encoding="utf-8", + errors="replace", + timeout=10, + check=False, + ) + if result.returncode != 0 and service.status() != settled_status: + detail = (result.stderr or result.stdout).strip() + raise RuntimeError(detail or f"sc.exe {verb} failed with {result.returncode}") - service = psutil.win_service_get(name) + +def _process_create_time(psutil, pid: int, label: str, when: str) -> float: + try: + return float(psutil.Process(int(pid)).create_time()) + except Exception as exc: + raise RuntimeError(f"Windows {label} process identity is unavailable {when}") from exc + + +def _verify_service_identities( + psutil, name: str, service, expected_service_identity, expected_gateway_identity +) -> None: + """Refuse to stop *name* unless SCM state, service/gateway process identities and ancestry all still match.""" if expected_service_identity is not None: try: current_status = str(service.status()) @@ -761,12 +746,7 @@ def _stop_windows_gateway_service( if identity is None: continue pid, create_time = identity - try: - current = float(psutil.Process(int(pid)).create_time()) - except Exception as exc: - raise RuntimeError( - f"Windows {label} process identity is unavailable before stop" - ) from exc + current = _process_create_time(psutil, pid, label, "before stop") if abs(current - float(create_time)) > 0.001: raise RuntimeError(f"Windows {label} process identity changed before stop") if expected_service_identity is not None and expected_gateway_identity is not None: @@ -780,18 +760,22 @@ def _stop_windows_gateway_service( ) from exc if service_pid not in ancestor_pids: raise RuntimeError(f"Windows gateway is no longer owned by service {name}") - result = subprocess.run( - ["sc.exe", "stop", name], - capture_output=True, - text=True, - encoding="utf-8", - errors="replace", - timeout=10, - check=False, - ) - if result.returncode != 0 and service.status() != "stopped": - detail = (result.stderr or result.stdout).strip() - raise RuntimeError(detail or f"sc.exe stop failed with {result.returncode}") + + +def _stop_windows_gateway_service( + name: str, + *, + expected_processes: tuple[tuple[int, float], ...] = (), + expected_service_identity: tuple[int, float] | None = None, + expected_gateway_identity: tuple[int, float] | None = None, + timeout: float = 30.0, +) -> None: + """Stop one verified Windows service and wait until SCM reports it down.""" + import psutil # noqa: PLC0415 + + service = psutil.win_service_get(name) + _verify_service_identities(psutil, name, service, expected_service_identity, expected_gateway_identity) + _sc_exe("stop", name, service, "stopped") def _original_process_is_alive(pid: int, create_time: float) -> bool: try: @@ -802,29 +786,17 @@ def _stop_windows_gateway_service( return True # AccessDenied/unknown: fail closed, venv may still be locked return abs(current - create_time) <= 0.001 - alive = [ - pid - for pid, create_time in expected_processes - if _original_process_is_alive(pid, create_time) - ] + def _alive() -> list[int]: + return [pid for pid, create_time in expected_processes if _original_process_is_alive(pid, create_time)] + deadline = _time.monotonic() + timeout while _time.monotonic() < deadline: - service_stopped = service.status() == "stopped" - alive = [ - pid - for pid, create_time in expected_processes - if _original_process_is_alive(pid, create_time) - ] - if service_stopped and not alive: + if service.status() == "stopped" and not _alive(): return _time.sleep(0.2) if service.status() == "stopped": # Lingering matching-identity processes make venv mutation unsafe — fail closed. - alive_after_stop = [ - pid - for pid, create_time in expected_processes - if _original_process_is_alive(pid, create_time) - ] + alive_after_stop = _alive() if alive_after_stop: raise RuntimeError( f"Windows service {name} stopped but its process tree is still alive: " @@ -841,18 +813,7 @@ def _start_windows_gateway_service(name: str, *, timeout: float = 30.0) -> None: import psutil # noqa: PLC0415 service = psutil.win_service_get(name) - result = subprocess.run( - ["sc.exe", "start", name], - capture_output=True, - text=True, - encoding="utf-8", - errors="replace", - timeout=10, - check=False, - ) - if result.returncode != 0 and service.status() != "running": - detail = (result.stderr or result.stdout).strip() - raise RuntimeError(detail or f"sc.exe start failed with {result.returncode}") + _sc_exe("start", name, service, "running") deadline = _time.monotonic() + timeout while _time.monotonic() < deadline: if service.status() == "running": @@ -884,15 +845,13 @@ def _restore_windows_gateway_service(name: str, *, timeout: float = 60.0) -> Non def _windows_cold_start_plan() -> dict | None: """Pause token for the no-running-gateway case: cold-start after update when an autostart entry exists. - Desktop-owned lifecycle -> ``None`` (spawning ``gateway run`` beside Desktop races ports/state). + An installed autostart entry is an explicit "I want a gateway" signal; a gateway that died between + updates would otherwise stay down until next login (resume only relaunches what was running). + Desktop-owned lifecycle -> ``None`` (spawning ``gateway run`` beside Desktop races ports/state); + the skip is ownership, not liveness. """ from hermes_cli.update_cmd import _desktop_owns_gateway_lifecycle - # No gateway running, but an installed autostart entry is an explicit "I want a - # gateway" signal; a gateway that died between updates would otherwise stay down - # until next login (resume only relaunches what was running). Cold-start after update. - # Exception: Desktop owns the lifecycle — spawning ``gateway run`` beside it races - # ports/state. The skip is ownership, not liveness. with _best_effort('Could not check Desktop gateway-lifecycle ownership before update: %s'): if _desktop_owns_gateway_lifecycle(): logger.debug( @@ -922,8 +881,6 @@ def _pause_windows_gateway_services(service_gateways, token: dict, profiles: dic """ from hermes_cli.update_cmd import _restore_windows_gateway_service, _stop_windows_gateway_service - # Stop SCM services only after every fallible ordinary-gateway step; from here any - # error restores attempted services and already-paused gateways before aborting. paused_services = [] current_service_name = None try: @@ -979,28 +936,13 @@ def _pause_windows_gateway_services(service_gateways, token: dict, profiles: dic raise RuntimeError(detail) from exc -def _pause_windows_gateways_for_update() -> dict | None: - """Stop running Windows gateways before mutating the checkout or venv. - - Scheduled/startup gateways run via pythonw.exe, invisible to the hermes.exe instance guard, yet keep files - locked during ``git``/``uv``. Stop only PIDs the gateway discovery code identifies. - """ - from hermes_cli.update_cmd import _m - - if not _m()._is_windows(): - return None - - try: - from gateway.status import get_process_start_time, terminate_pid - from hermes_cli.gateway import ( - _capture_gateway_argv, - _get_restart_drain_timeout, - find_gateway_pids, - find_profile_gateway_processes, - find_windows_gateway_services, - ) - except Exception as exc: - raise RuntimeError(f"Could not prepare Windows gateway pause for update: {exc}") from exc +def _discover_windows_gateways(): + """``(profile_processes, service_gateways, running_pids)`` for the pause; any indeterminate probe aborts.""" + from hermes_cli.gateway import ( + find_gateway_pids, + find_profile_gateway_processes, + find_windows_gateway_services, + ) try: profile_process_list = find_profile_gateway_processes(strict=True) @@ -1026,9 +968,15 @@ def _pause_windows_gateways_for_update() -> dict | None: ) except Exception as exc: raise RuntimeError(f"Could not discover Windows gateway PIDs before update: {exc}") from exc - if not running_pids: - return _windows_cold_start_plan() + return profile_processes, service_gateways, service_gateway_pids, running_pids + +def _request_socket_pauses(running_pids, profile_processes, service_gateway_pids): + """Marker + socket-first pause for every profile-mapped gateway; ``(profiles, mapped_pids, socket_acks)``. + + Socket ACK = the gateway drains and exits by its own graceful path. No answer (older + gateway) -> the marker poll / force-kill ladder in the caller. + """ profiles: dict[str, int] = {} mapped_pids = [] socket_acks: list[dict] = [] @@ -1041,8 +989,6 @@ def _pause_windows_gateways_for_update() -> dict | None: profiles[str(proc.profile)] = int(pid) mapped_pids.append(int(pid)) _write_update_planned_stop_marker(Path(proc.path), int(pid)) - # Socket-first pause: ask the gateway to drain and exit itself (ACK = its own - # graceful path). No answer (older gateway) -> marker poll / force-kill ladder below. try: from gateway.control_socket import pause_gateway_for_update @@ -1051,20 +997,18 @@ def _pause_windows_gateways_for_update() -> dict | None: socket_acks.append(ack) except Exception as exc: logger.debug("Socket pause unavailable for gateway %s: %s", pid, exc) + return profiles, mapped_pids, socket_acks - # Resolve venv-side launchers BEFORE draining: a dead worker's parent cannot be - # recovered (NoSuchProcess). The launcher keeps ``.pyd`` mapped and would trip the - # venv-holder guard after the gateway stopped; it is killed with the survivors. - launcher_pids = _m()._venv_launcher_ancestors(mapped_pids) - print("→ Stopping Windows gateway process(es) before updating Hermes...") +def _gateway_drain_timeout(socket_acks: list[dict]) -> float: + """Drain budget: configured restart drain (>= 1s), raised to a socket-paused gateway's declared + ACTIVE-TURN budget + teardown grace so it isn't force-killed mid-turn.""" + from hermes_cli.gateway import _get_restart_drain_timeout try: drain_timeout = max(float(_get_restart_drain_timeout()), 1.0) except Exception: drain_timeout = 10.0 if socket_acks: - # A socket-paused gateway drains its ACTIVE TURN first; honor its declared - # budget (+ teardown grace) so it isn't force-killed mid-turn. with suppress(Exception): declared = max(float(a.get("drain_timeout") or 0.0) for a in socket_acks) drain_timeout = max(drain_timeout, declared + 10.0) @@ -1072,6 +1016,41 @@ def _pause_windows_gateways_for_update() -> dict | None: f" → {len(socket_acks)} gateway(s) ACKed socket pause; " f"waiting up to {int(drain_timeout)}s for graceful exit" ) + return drain_timeout + + +def _pause_windows_gateways_for_update() -> dict | None: + """Stop running Windows gateways before mutating the checkout or venv. + + Scheduled/startup gateways run via pythonw.exe, invisible to the hermes.exe instance guard, yet keep files + locked during ``git``/``uv``. Stop only PIDs the gateway discovery code identifies. + """ + from hermes_cli.update_cmd import _m + + if not _m()._is_windows(): + return None + + try: + from gateway.status import get_process_start_time, terminate_pid + from hermes_cli.gateway import _capture_gateway_argv + except Exception as exc: + raise RuntimeError(f"Could not prepare Windows gateway pause for update: {exc}") from exc + + profile_processes, service_gateways, service_gateway_pids, running_pids = _discover_windows_gateways() + if not running_pids: + return _windows_cold_start_plan() + + profiles, mapped_pids, socket_acks = _request_socket_pauses( + running_pids, profile_processes, service_gateway_pids + ) + + # Resolve venv-side launchers BEFORE draining: a dead worker's parent cannot be + # recovered (NoSuchProcess). The launcher keeps ``.pyd`` mapped and would trip the + # venv-holder guard after the gateway stopped; it is killed with the survivors. + launcher_pids = _m()._venv_launcher_ancestors(mapped_pids) + + print("→ Stopping Windows gateway process(es) before updating Hermes...") + drain_timeout = _gateway_drain_timeout(socket_acks) survivors = _m()._wait_for_windows_update_gateway_exit(mapped_pids, timeout=drain_timeout) unmapped_pids = [ pid @@ -1105,9 +1084,8 @@ def _pause_windows_gateways_for_update() -> dict | None: print(f" → Force-stopped {len(force_killed)} gateway process(es)") if unmapped_pids: - respawnable = sum(1 for u in unmapped if u.get("argv")) print(f" → Stopped {len(unmapped_pids)} gateway process(es) without profile mapping") - if respawnable < len(unmapped_pids): + if sum(1 for u in unmapped if u.get("argv")) < len(unmapped_pids): # No recoverable cmdline (psutil missing, access denied, gone): manual restart. print(" Restart manually after update: hermes gateway run") @@ -1241,19 +1219,9 @@ def _refresh_bootstrap_cache_scripts(branch: str = "main") -> None: ) -def _resume_windows_gateways_after_update(token: dict | None) -> None: - """Restart Windows profile gateways previously paused for update.""" - from hermes_cli.update_cmd import _m, _start_windows_gateway_service - if not token or not token.get("resume_needed"): - return - if not _m()._is_windows(): - token["resume_needed"] = False - return - - # Regenerate launcher scripts before respawning so a legacy pythonw-era - # autostart entry comes back on the current design at next login too. - _m()._refresh_windows_gateway_launchers() - +def _resume_windows_services(token: dict) -> None: + """Restart the SCM services recorded on *token*; failed ones stay on the token so a retry sees them.""" + from hermes_cli.update_cmd import _start_windows_gateway_service services = list(token.get("services") or []) token.setdefault("expected_services", list(services)) verified_restarts = list(token.get("restarted_services") or []) @@ -1274,30 +1242,25 @@ def _resume_windows_gateways_after_update(token: dict | None) -> None: print(f" ⚠ Could not restart Windows gateway service: {service_name}") failed_services.append(str(service_name)) + token["restarted_services"] = verified_restarts if failed_services: token["services"] = failed_services - token["restarted_services"] = verified_restarts raise RuntimeError( "Could not restart Windows gateway service(s): " + ", ".join(failed_services) ) token["services"] = [] - token["restarted_services"] = verified_restarts if restarted_services: print() print(" ✓ Restarted Windows gateway service(s): " + ", ".join(restarted_services)) - profiles = token.get("profiles") or {} - unmapped = token.get("unmapped") or [] - cold_start = bool(token.get("cold_start_if_installed")) - if not profiles and not any(u.get("argv") for u in unmapped): - if cold_start: - if not _m()._cold_start_windows_gateway_after_update(): - raise RuntimeError("Windows gateway cold-start was not verified") - token["cold_start_if_installed"] = False - token["resume_needed"] = False - return +def _relaunch_paused_gateways(token: dict, profiles: dict, unmapped: list) -> tuple[list[str], int]: + """Relaunch profile gateways and replay unmapped argv; ``(relaunched_profiles, unmapped_count)``. + + Failed relaunches stay on the token (and off ``relaunched_profiles``) so plan-vs-execution + reconciliation still surfaces them — Windows has no watcher to recover them. + """ try: from hermes_cli.gateway import ( launch_detached_gateway_restart_by_cmdline, @@ -1321,13 +1284,8 @@ def _resume_windows_gateways_after_update(token: dict | None) -> None: exc, ) failed_profiles[str(profile)] = int(old_pid) - - # Feed the plan-vs-execution reconciliation (else a relaunched gateway is reported - # "unaccounted", exit 1). Failed relaunches are deliberately left off so they - # still surface (Windows has no watcher to recover them). token["relaunched_profiles"] = relaunched - # Respawn unmapped gateways by replaying the argv snapshotted before the kill. unmapped_relaunched = 0 failed_unmapped = [] for entry in unmapped: @@ -1353,35 +1311,66 @@ def _resume_windows_gateways_after_update(token: dict | None) -> None: token["unmapped"] = failed_unmapped if failed_profiles or failed_unmapped: raise RuntimeError("Could not restart every paused Windows gateway") + return relaunched, unmapped_relaunched - # A truthy launch only proves the watcher was created; a parent Job Object denying - # CREATE_BREAKAWAY_FROM_JOB can kill the gateway on updater teardown. Verify with - # the same liveness poll every spawn path uses; all_profiles=True covers the fleet. + +def _verify_relaunched_gateways_alive(token: dict, profiles: dict, unmapped: list) -> None: + """Gate success on the shared liveness poll: a truthy launch only proves the watcher was created. + + A parent Job Object denying CREATE_BREAKAWAY_FROM_JOB can kill the gateway on updater teardown; + ``all_profiles=True`` covers the fleet. Vouched PIDs are persisted so a death AFTER updater exit + is reported by the next CLI invocation (best-effort). + """ + try: + from hermes_cli import gateway_windows + except Exception as exc: + raise RuntimeError(f"Could not load Windows gateway liveness helpers: {exc}") from exc + ready_pids = gateway_windows._wait_for_gateway_ready(timeout_s=30.0, all_profiles=True) + if not ready_pids: + token["profiles"] = dict(profiles) + token["unmapped"] = list(unmapped) + print() + print( + " ⚠ Windows gateway restart could not be verified — no stable " + "gateway process appeared after relaunch." + ) + print( + " (The respawned gateway may have been killed by a parent " + "Job Object during updater teardown, #48820.)" + ) + print(" Recover with: hermes gateway restart") + raise RuntimeError("Windows gateway relaunch after update was not verified alive") + with suppress(Exception): + gateway_windows._write_start_attestation(ready_pids, "post-update relaunch") + + +def _resume_windows_gateways_after_update(token: dict | None) -> None: + """Restart Windows profile gateways previously paused for update.""" + from hermes_cli.update_cmd import _m + if not token or not token.get("resume_needed"): + return + if not _m()._is_windows(): + token["resume_needed"] = False + return + + # Regenerate launcher scripts before respawning so a legacy pythonw-era + # autostart entry comes back on the current design at next login too. + _m()._refresh_windows_gateway_launchers() + _resume_windows_services(token) + + profiles = token.get("profiles") or {} + unmapped = token.get("unmapped") or [] + if not profiles and not any(u.get("argv") for u in unmapped): + if token.get("cold_start_if_installed"): + if not _m()._cold_start_windows_gateway_after_update(): + raise RuntimeError("Windows gateway cold-start was not verified") + token["cold_start_if_installed"] = False + token["resume_needed"] = False + return + + relaunched, unmapped_relaunched = _relaunch_paused_gateways(token, profiles, unmapped) if relaunched or unmapped_relaunched: - try: - from hermes_cli import gateway_windows - except Exception as exc: - raise RuntimeError(f"Could not load Windows gateway liveness helpers: {exc}") from exc - ready_pids = gateway_windows._wait_for_gateway_ready(timeout_s=30.0, all_profiles=True) - if not ready_pids: - token["profiles"] = dict(profiles) - token["unmapped"] = list(unmapped) - print() - print( - " ⚠ Windows gateway restart could not be verified — no stable " - "gateway process appeared after relaunch." - ) - print( - " (The respawned gateway may have been killed by a parent " - "Job Object during updater teardown, #48820.)" - ) - print(" Recover with: hermes gateway restart") - raise RuntimeError("Windows gateway relaunch after update was not verified alive") - # Persist vouched PIDs so a death AFTER updater exit is reported by the - # next CLI invocation. Best-effort. - with suppress(Exception): - gateway_windows._write_start_attestation(ready_pids, "post-update relaunch") - + _verify_relaunched_gateways_alive(token, profiles, unmapped) token["resume_needed"] = False if relaunched: @@ -1408,25 +1397,25 @@ def _resume_windows_gateways_and_merge_outcome(outcome, _windows_gateway_resume, _write_gateway_update_exit_code(False) if isinstance(_windows_gateway_resume, dict): + def _extend_unique(target: list, items) -> None: + for item in items: + if item not in target: + target.append(item) + # Failed relaunches are absent from the token so they still surface. Best-effort. with _best_effort('Could not merge Windows relaunch outcome into fleet reconciliation bookkeeping: %s'): - for _win_profile in _windows_gateway_resume.get("relaunched_profiles") or []: - if _win_profile not in outcome.relaunched_profiles: - outcome.relaunched_profiles.append(_win_profile) + _extend_unique(outcome.relaunched_profiles, _windows_gateway_resume.get("relaunched_profiles") or []) windows_restarted = list(_windows_gateway_resume.get("restarted_services") or []) - for service_name in windows_restarted: - if service_name not in outcome.restarted_services: - outcome.restarted_services.append(service_name) service_profiles = _windows_gateway_resume.get("service_profiles") or {} - for service_name in windows_restarted: - profile_name = service_profiles.get(service_name) - if profile_name and profile_name not in outcome.relaunched_profiles: - outcome.relaunched_profiles.append(profile_name) - pending_services = list(_windows_gateway_resume.get("services") or []) - for service_name in pending_services: - label = str(service_profiles.get(service_name) or service_name) - if label not in outcome.failed_or_stale_units: - outcome.failed_or_stale_units.append(label) + _extend_unique(outcome.restarted_services, windows_restarted) + _extend_unique( + outcome.relaunched_profiles, + (p for p in (service_profiles.get(n) for n in windows_restarted) if p), + ) + _extend_unique( + outcome.failed_or_stale_units, + (str(service_profiles.get(n) or n) for n in (_windows_gateway_resume.get("services") or [])), + ) with suppress(Exception): from hermes_cli.update_receipt import record_gateway_restart @@ -1445,6 +1434,46 @@ def _resume_windows_gateways_and_merge_outcome(outcome, _windows_gateway_resume, ) +def _reap_and_rescan(message: str, pids, stop=None) -> list[tuple[int, str, str]]: + """Announce *message*, stop *pids* (tree-kill unless *stop* given), settle 1s, re-scan venv holders.""" + from hermes_cli.update_cmd import _m + print(message) + (stop or _m()._stop_process_trees)(pids) + _time.sleep(1.0) + return _m()._detect_venv_python_processes() + + +def _terminate_leftover_gateways(pids) -> None: + """Force-stop leftover gateways one by one; a failure is logged, never raised.""" + from gateway.status import get_process_start_time, terminate_pid + + for _pid in pids: + try: + pid_int = int(_pid) + terminate_pid(pid_int, force=True, expected_start_time=get_process_start_time(pid_int)) + except Exception as exc: + logger.debug("Could not stop leftover gateway %s: %s", _pid, exc) + + +def _in_handoff_without_live_shim(args) -> bool: + """GUI hand-off gate: ``--gateway`` + update-incomplete marker AND no live ``hermes.exe`` shim. + + Fail closed: unverifiable marker or shim state reads as "not a hand-off" / "live shim". + """ + from hermes_cli.update_cmd import _m + try: + handoff = bool(getattr(args, "gateway", False)) and _m()._update_marker_path().exists() + except Exception: + return False + if not handoff: + return False + try: + scripts_dir = _m()._venv_scripts_dir() + return scripts_dir is not None and not _m()._detect_concurrent_hermes_instances(scripts_dir) + except Exception: + return False + + def _clear_windows_venv_holders_or_exit(args, gateway_mode: bool, _windows_gateway_resume): """Windows: stop every venv-python holder we can positively identify, else resume paused gateways and exit 2. @@ -1455,104 +1484,73 @@ def _clear_windows_venv_holders_or_exit(args, gateway_mode: bool, _windows_gatew from hermes_cli.update_cmd import _m, _record_update_step, _refuse_gateway_ancestor_tree_kill _venv_holders = _m()._detect_venv_python_processes() if _venv_holders: + # Gateways the pause machinery owns (respawned in the pause->guard window or + # unmapped spawn path): stop and re-check; post-update resume brings them back. _gateway_holders = _m()._leftover_pausable_gateway_pids(_venv_holders) if _gateway_holders is not None: - if _refuse_gateway_ancestor_tree_kill( - _gateway_holders, gateway_mode=gateway_mode - ): + if _refuse_gateway_ancestor_tree_kill(_gateway_holders, gateway_mode=gateway_mode): _m()._resume_windows_gateways_after_update(_windows_gateway_resume) sys.exit(2) - # Gateways the pause machinery owns (respawned in the pause->guard window or - # unmapped spawn path): stop and re-check; post-update resume brings them back. - from gateway.status import get_process_start_time, terminate_pid - - print( + _venv_holders = _reap_and_rescan( f" ⚠ {len(_gateway_holders)} gateway process(es) still " - "hold the venv after the pause; stopping them" + "hold the venv after the pause; stopping them", + _gateway_holders, + stop=_terminate_leftover_gateways, ) - for _pid in _gateway_holders: - try: - pid_int = int(_pid) - terminate_pid( - pid_int, - force=True, - expected_start_time=get_process_start_time(pid_int), - ) - except Exception as exc: - logger.debug("Could not stop leftover gateway %s: %s", _pid, exc) - _time.sleep(1.0) - _venv_holders = _m()._detect_venv_python_processes() if _venv_holders: # Positive-identity rung (any context): spawn ledger proves the holder is an # orphaned backend (self-registered, spawner provably dead). No PPID archaeology. _ledger_backends = _m()._ledger_reapable_backend_pids(_venv_holders) if _ledger_backends: - print( + _venv_holders = _reap_and_rescan( f" ⚠ {len(_ledger_backends)} ledger-identified orphaned " - "Hermes backend process(es) hold the venv; stopping their trees" + "Hermes backend process(es) hold the venv; stopping their trees", + _ledger_backends, ) - _m()._stop_process_trees(_ledger_backends) - _time.sleep(1.0) - _venv_holders = _m()._detect_venv_python_processes() if _venv_holders: + # Desktop `serve` backends whose app is GONE: nothing respawns an orphan, so + # reap the tree. Live-Desktop backends return None and keep the refusal. _orphan_backends = _m()._orphaned_desktop_backend_pids(_venv_holders) if _orphan_backends: - # Desktop `serve` backends whose app is GONE: nothing respawns an orphan, so - # reap the tree. Live-Desktop backends return None and keep the refusal. - print( + _venv_holders = _reap_and_rescan( f" ⚠ {len(_orphan_backends)} orphaned Desktop backend " - "process(es) still hold the venv; stopping their trees" + "process(es) still hold the venv; stopping their trees", + _orphan_backends, ) - _m()._stop_process_trees(_orphan_backends) - _time.sleep(1.0) - _venv_holders = _m()._detect_venv_python_processes() if _venv_holders: # Manual serve/dashboard rung (e.g. `hermes serve --host ` for a REMOTE Desktop): # ledger identity only (spawner dead; Desktop-owned keep the refusal). Stop and # register an idempotent atexit relaunch on the SAME host/port/profile — success or failure. _serve_entries = _m()._ledger_manual_serve_holders(_venv_holders) if _serve_entries: - print( + def _stop_and_park(pids): + _m()._stop_process_trees(pids) + _record_update_step("serve_pause", True, f"stopped={len(_serve_entries)}") + import atexit as _serve_atexit + + _serve_atexit.register( + _m()._relaunch_stopped_serves, {"pending": True, "entries": _serve_entries} + ) + + _venv_holders = _reap_and_rescan( f" ⚠ {len(_serve_entries)} manual serve/dashboard " "backend(s) hold the venv; stopping them for the update " - "(they will be relaunched on their recorded endpoints)" + "(they will be relaunched on their recorded endpoints)", + [int(e["pid"]) for e in _serve_entries], + stop=_stop_and_park, + ) + if _venv_holders and _in_handoff_without_live_shim(args): + # Final rung: in a GUI hand-off the Desktop is contractually gone; surviving `serve` + # backends are leaks even with a live parent (which made the orphan-only rung bail + # and hang) — reap by cmdline. + _handoff_backends = _m()._handoff_reapable_backend_pids(_venv_holders) + if _handoff_backends: + _venv_holders = _reap_and_rescan( + f" ⚠ {len(_handoff_backends)} Hermes backend process(es) " + "still hold the venv after the Desktop hand-off; " + "stopping their trees", + _handoff_backends, ) - _m()._stop_process_trees([int(e["pid"]) for e in _serve_entries]) - _serve_resume_token = {"pending": True, "entries": _serve_entries} - _record_update_step("serve_pause", True, f"stopped={len(_serve_entries)}") - import atexit as _serve_atexit - - _serve_atexit.register(_m()._relaunch_stopped_serves, _serve_resume_token) - _time.sleep(1.0) - _venv_holders = _m()._detect_venv_python_processes() - if _venv_holders: - # Final rung: in a GUI hand-off (`--gateway` + update-incomplete marker) the Desktop - # is contractually gone; surviving `serve` backends are leaks even with a live - # parent (which made the orphan-only rung bail and hang) — reap by cmdline. - _handoff = False - try: - _handoff = bool(getattr(args, "gateway", False)) and _m()._update_marker_path().exists() - except Exception: - _handoff = False - # Fail closed: unverifiable shim state is treated as a live shim (keep refusing). - _no_live_shim = False - try: - _scripts_dir = _m()._venv_scripts_dir() - if _scripts_dir is not None: - _no_live_shim = not _m()._detect_concurrent_hermes_instances(_scripts_dir) - except Exception: - _no_live_shim = False - if _handoff and _no_live_shim: - _handoff_backends = _m()._handoff_reapable_backend_pids(_venv_holders) - if _handoff_backends: - print( - f" ⚠ {len(_handoff_backends)} Hermes backend process(es) " - "still hold the venv after the Desktop hand-off; " - "stopping their trees" - ) - _m()._stop_process_trees(_handoff_backends) - _time.sleep(1.0) - _venv_holders = _m()._detect_venv_python_processes() if _venv_holders: print(_format_venv_python_holders_message(_venv_holders)) _m()._resume_windows_gateways_after_update(_windows_gateway_resume) diff --git a/hermes_cli/update_restart_recovery.py b/hermes_cli/update_restart_recovery.py index 3e001f9640..3b92865c69 100644 --- a/hermes_cli/update_restart_recovery.py +++ b/hermes_cli/update_restart_recovery.py @@ -110,11 +110,7 @@ def _child_environment() -> dict[str, str]: return env -def _run_profile_restart( - profile: str, - *, - run: Callable[..., Any], -) -> bool: +def _run_profile_restart(profile: str, *, run: Callable[..., Any]) -> bool: """Run one profile restart without inheriting the updater's process state.""" kwargs: dict[str, Any] = {"stdin": subprocess.DEVNULL, "env": _child_environment()} if os.name == "nt": @@ -157,27 +153,22 @@ def restart_profiles( supervisors: Mapping[str, str] | None = None, run: Callable[..., Any] = subprocess.run, ) -> dict[str, list[str]]: - """Restart the supplied profiles and return per-profile terminal results. - - The caller supplies only profiles whose inventory identified a service supervisor. + """Restart the supplied profiles (only ones whose inventory identified a service supervisor). A profile only lands in ``verified`` when its supervisor is systemd and ``systemctl --user is- active`` independently confirms the unit after the relaunch command succeeded. """ supervisors = supervisors or {} - normalized = sorted({profile for profile in profiles if isinstance(profile, str) and profile}) - verified: list[str] = [] - relaunch_attempted: list[str] = [] - failed: list[str] = [] - for profile in normalized: + result: dict[str, list[str]] = {"verified": [], "relaunch_attempted": [], "failed": []} + for profile in sorted({p for p in profiles if isinstance(p, str) and p}): if not _run_profile_restart(profile, run=run): - failed.append(profile) - continue - if supervisors.get(profile) == "systemd" and _systemd_verified_active(profile, run=run): - verified.append(profile) + bucket = "failed" + elif supervisors.get(profile) == "systemd" and _systemd_verified_active(profile, run=run): + bucket = "verified" else: - relaunch_attempted.append(profile) - return {"verified": verified, "relaunch_attempted": relaunch_attempted, "failed": failed} + bucket = "relaunch_attempted" + result[bucket].append(profile) + return result def _systemctl_scopes() -> list[tuple[str, list[str]]]: @@ -206,10 +197,8 @@ def _listed_serve_units(scope: list[str], *, run: Callable[..., Any]) -> list[st ) if result is None: return [] - # The glob is a systemd pattern, not a name gate: `hermes-serve*` also - # matches the unrelated `hermes-server.service`. Require the exact - # base unit or the hyphenated profile family, same shape as the - # in-process phase's own name gate. + # The glob is a systemd pattern, not a name gate (`hermes-serve*` also matches + # `hermes-server.service`): require the exact base unit or the profile family. units: list[str] = [] for line in (getattr(result, "stdout", "") or "").splitlines(): parts = line.split() @@ -309,9 +298,8 @@ def restart_serve_units( unit and therefore cannot be touched here. """ skipped_qualified, skipped_legacy = _normalized_skips(skip_units) - # (scope, base unit) -> replaced? A unit name can exist in BOTH the user - # and the system scope; each is a separate process and each is proven, - # reported and accounted for on its own. + # (scope, base unit) -> replaced? The same unit name in the user and the system + # scope is two processes; each is proven, reported and accounted for on its own. outcomes: dict[tuple[str, str], bool] = {} seen: set[tuple[str, str]] = set() for scope_label, scope in _systemctl_scopes(): @@ -322,26 +310,15 @@ def restart_serve_units( continue seen.add(target) if not _unit_is_active(scope, unit, run=run): - # Not running: nothing is serving a stale generation from it. - continue + continue # not running: nothing serves a stale generation from it previous_pid = _unit_main_pid(scope, unit, run=run) - if previous_pid <= 0: - # Active with no readable main process: a replacement cannot - # be observed, so it cannot be claimed. Restarting blind and - # reporting success is the failure mode this module exists to - # remove. - outcomes[target] = False - continue - result = _run_quiet( - run, scope + ["--no-ask-password", "restart", unit], timeout=_UNIT_RESTART_TIMEOUT - ) - if not _succeeded(result): - # Includes the unprivileged system-scope case. We do not probe - # for sudo here: an unverifiable unit must read as failed so - # the update stays explicitly incomplete. - outcomes[target] = False - continue - outcomes[target] = _serve_unit_replaced(scope, unit, previous_pid, run=run, sleep=sleep) + # No readable main PID: a replacement can't be observed, so it can't be claimed + # (restarting blind and reporting success is the failure mode this module removes). + # A failed restart includes the unprivileged system-scope case; no sudo probe — + # an unverifiable unit must read as failed so the update stays explicitly incomplete. + outcomes[target] = previous_pid > 0 and _succeeded( + _run_quiet(run, scope + ["--no-ask-password", "restart", unit], timeout=_UNIT_RESTART_TIMEOUT) + ) and _serve_unit_replaced(scope, unit, previous_pid, run=run, sleep=sleep) # ``/`` is the only identity this module reports for a serve unit. return { "verified": sorted(f"{scope}/{base}" for (scope, base), ok in outcomes.items() if ok),