From e64942a0184ccfd54e4fbbb95b370e751c502d6e Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 21:18:02 -0700 Subject: [PATCH] refactor(hermes_cli/update): compact fleet restart bookkeeping into outcome dataclass; split maint backup/notice phases --- hermes_cli/update_cmd_fleet.py | 942 +++++++----------- hermes_cli/update_cmd_maint.py | 518 +++++----- .../test_update_fleet_check_fail_closed.py | 10 +- 3 files changed, 619 insertions(+), 851 deletions(-) diff --git a/hermes_cli/update_cmd_fleet.py b/hermes_cli/update_cmd_fleet.py index 372f7765de..c1ce2c152a 100644 --- a/hermes_cli/update_cmd_fleet.py +++ b/hermes_cli/update_cmd_fleet.py @@ -19,18 +19,21 @@ from hermes_cli.update_cmd_common import _best_effort # Log-record parity with the origin module. logger = logging.getLogger("hermes_cli.update_cmd") - -def _write_gateway_update_exit_code(ok: bool) -> None: - from hermes_cli.update_cmd import get_hermes_home - path = get_hermes_home() / ".update_exit_code" - with suppress(OSError): - path.write_text("0" if ok else "1", encoding="utf-8") - - # Under HERMES_HOME (not next to the venv): records the fleet-restart obligation # after a pull advanced HEAD; cleared only when the restart completes or nothing ran. _FLEET_RESTART_PENDING_NAME = "fleet_restart_pending" +_FRESH_RESTART_SUPERVISORS = frozenset({"systemd", "launchd", "service", "s6"}) + +_SYSTEMD_SCOPES = (("user", ["systemctl", "--user"]), ("system", ["systemctl"])) +_LIST_GATEWAY_UNITS = ["list-units", "hermes-gateway*", "hermes-serve*", "--plain", "--no-legend", "--no-pager"] + + +def _write_gateway_update_exit_code(ok: bool) -> None: + from hermes_cli.update_cmd import get_hermes_home + with suppress(OSError): + (get_hermes_home() / ".update_exit_code").write_text("0" if ok else "1", encoding="utf-8") + def _fleet_restart_pending_marker_path() -> Path: """HERMES_HOME breadcrumb for a pull that has not yet restarted the fleet.""" @@ -74,18 +77,13 @@ def _current_checkout_sha() -> str | None: def _receipt_looks_unfinished(receipt: dict) -> bool: """True when *receipt* is from an update that did not finish cleanly.""" - if receipt.get("stop_reason"): - return True - exit_code = receipt.get("exit_code") - if exit_code not in (0, None): - return True - outcome = receipt.get("outcome") - if outcome in ("failed", "partial", "running"): - return True gateway_restart = receipt.get("gateway_restart") - if isinstance(gateway_restart, dict) and gateway_restart.get("incomplete"): - return True - return False + return bool( + receipt.get("stop_reason") + or receipt.get("exit_code") not in (0, None) + or receipt.get("outcome") in ("failed", "partial", "running") + or (isinstance(gateway_restart, dict) and gateway_restart.get("incomplete")) + ) def _receipt_reports_stale_runtime(expected_sha: str | None = None) -> bool: @@ -104,8 +102,7 @@ def _receipt_reports_stale_runtime(expected_sha: str | None = None) -> bool: receipt = None if not isinstance(receipt, dict): return False - if not expected_sha: - expected_sha = _current_checkout_sha() + expected_sha = expected_sha or _current_checkout_sha() if not expected_sha: return False @@ -114,24 +111,21 @@ def _receipt_reports_stale_runtime(expected_sha: str | None = None) -> bool: fleet = receipt.get("fleet") if isinstance(fleet, list) and fleet: - for entry in fleet: - if not isinstance(entry, dict): - continue - if entry.get("state") == "stale": - return True - if _sha_mismatch(entry.get("code_sha")): - return True - return False + return any( + isinstance(entry, dict) + and (entry.get("state") == "stale" or _sha_mismatch(entry.get("code_sha"))) + for entry in fleet + ) if not _receipt_looks_unfinished(receipt): return False plan = receipt.get("plan") if not isinstance(plan, dict): return False - for runtime in plan.get("runtimes") or []: - if isinstance(runtime, dict) and _sha_mismatch(runtime.get("code_sha")): - return True - return False + return any( + isinstance(runtime, dict) and _sha_mismatch(runtime.get("code_sha")) + for runtime in plan.get("runtimes") or [] + ) def _pending_fleet_restart_needed() -> bool: @@ -158,45 +152,52 @@ def _warn_pending_fleet_restart(*, startup: bool = False) -> None: def _warn_pending_fleet_restart_on_startup() -> None: """Cheap CLI-startup hint. Never restarts; never raises.""" with suppress(Exception): - if not _pending_fleet_restart_needed(): - return - _warn_pending_fleet_restart(startup=True) + if _pending_fleet_restart_needed(): + _warn_pending_fleet_restart(startup=True) + + +def _systemd_gateway_unit_listings(on_list_timeout=None): + """Yield ``(scope, scope_cmd, list-units CompletedProcess)`` per systemd scope that answered. + + A missing systemctl skips the scope silently; a listing timeout skips it after + ``on_list_timeout(scope, exc)`` (when given) so the other scope is still processed. + """ + for scope, scope_cmd in _SYSTEMD_SCOPES: + try: + result = _systemctl(scope_cmd + _LIST_GATEWAY_UNITS, timeout=10) + except FileNotFoundError: + continue + except subprocess.TimeoutExpired as exc: + if on_list_timeout is not None: + on_list_timeout(scope, exc) + continue + yield scope, scope_cmd, result + + +def _needs_sudo(scope: str) -> bool: + return ( + scope == "system" + and hasattr(os, "geteuid") + and os.geteuid() != 0 # windows-footgun: ok — systemd path, Linux-only + ) def _restart_systemd_gateway_units_best_effort(failed: list) -> None: """Best-effort ``systemctl restart`` of every hermes-gateway/serve unit.""" - for scope, scope_cmd in ( - ("user", ["systemctl", "--user"]), - ("system", ["systemctl"]), - ): - try: - result = _systemctl( - scope_cmd + ["list-units", "hermes-gateway*", "hermes-serve*", - "--plain", "--no-legend", "--no-pager"], - timeout=10, - ) - except (FileNotFoundError, subprocess.TimeoutExpired): - continue + for scope, scope_cmd, result in _systemd_gateway_unit_listings(): if result.returncode != 0: continue def process_unit(svc_name: str, _scope=scope, _cmd=scope_cmd) -> None: restart_cmd = list(_cmd) + ["--no-ask-password", "restart", svc_name] - if ( - _scope == "system" - and hasattr(os, "geteuid") - and os.geteuid() != 0 # windows-footgun: ok — systemd path, Linux-only - ): + if _needs_sudo(_scope): restart_cmd = ["sudo", "-n"] + restart_cmd _systemctl(restart_cmd, timeout=30) - def on_timeout(svc_name: str, exc: subprocess.TimeoutExpired) -> None: - failed.append(svc_name) - _for_each_systemd_gateway_unit( result.stdout, process_unit=process_unit, - on_unit_timeout=on_timeout, + on_unit_timeout=lambda svc_name, exc: failed.append(svc_name), ) @@ -237,9 +238,8 @@ def _run_pending_fleet_restart() -> bool: if supports_systemd_services(): _restart_systemd_gateway_units_best_effort(failed) if is_macos(): - restarted: list = [] try: - _restart_macos_launchd_gateways(restarted, failed, 45.0) + _restart_macos_launchd_gateways([], failed, 45.0) except Exception as exc: logger.debug("Pending fleet restart: launchd failed: %s", exc) failed.append("launchd") @@ -252,7 +252,6 @@ def _run_pending_fleet_restart() -> bool: except Exception as exc: logger.debug("Pending fleet restart: Windows failed: %s", exc) failed.append("windows-gateway") - leftover: list = [] try: leftover = list(find_gateway_pids(all_profiles=True)) except Exception: @@ -267,7 +266,6 @@ def _run_pending_fleet_restart() -> bool: print(" ✓ Pending fleet restart completed.") return True except Exception as exc: - surviving = None try: surviving = list(find_gateway_pids(all_profiles=True)) except Exception: @@ -312,6 +310,17 @@ def _systemctl_reset_and_restart(manage_cmd: list, svc_name: str): return _systemctl(manage_cmd + ["restart", svc_name], timeout=15) +def _is_hermes_gateway_unit(unit: str) -> bool: + """Exact base unit or hyphenated profile family only: ``startswith("hermes-serve")`` + would accept ``hermes-server.service``.""" + return ( + unit == "hermes-gateway.service" + or unit.startswith("hermes-gateway-") + or unit == "hermes-serve.service" + or unit.startswith("hermes-serve-") + ) + + def _for_each_systemd_gateway_unit( list_units_stdout: str, *, @@ -328,16 +337,7 @@ def _for_each_systemd_gateway_unit( if not parts: continue unit = parts[0] - if not unit.endswith(".service"): - continue - # Name gate against stray lines. Exact base unit or hyphenated profile family - # only: ``startswith("hermes-serve")`` would accept ``hermes-server.service``. - if not ( - unit == "hermes-gateway.service" - or unit.startswith("hermes-gateway-") - or unit == "hermes-serve.service" - or unit.startswith("hermes-serve-") - ): + if not unit.endswith(".service") or not _is_hermes_gateway_unit(unit): continue svc_name = unit.removesuffix(".service") try: @@ -363,14 +363,7 @@ def _warn_incomplete_gateway_fleet_restart(failed_units: list) -> None: if not failed_units: return - # Preserve discovery order while de-duplicating. - seen = set() - ordered = [] - for name in failed_units: - if name in seen: - continue - seen.add(name) - ordered.append(name) + ordered = list(dict.fromkeys(failed_units)) # de-dup, discovery order print() print("⚠ Update incomplete — some units were not restarted:") for name in ordered: @@ -400,13 +393,11 @@ def _restart_launchd_gateway_after_update( """Restart the invoking profile's launchd gateway after an update. No ``launchctl list`` gating: a booted-out job (plist present, definition - deregistered) fails it, and it is session-scoped and can exit non-zero while the - job is alive — gating on it silently skipped the restart yet printed "Update - complete!". When the plist exists ``launchd_restart()`` always runs (drain, -k - kickstart, bootout/bootstrap ladder); every failure path is loud with a manual - recovery command. Returns ``(restarted_labels, failed_labels)``; with - ``supervision_verify`` success also requires a fresh supervised PID ("the call - returned" is not "supervised"). + deregistered) fails it, and it can exit non-zero while the job is alive — gating on it + silently skipped the restart yet printed "Update complete!". When the plist exists + ``launchd_restart()`` always runs; every failure path is loud with a manual recovery + command. Returns ``(restarted_labels, failed_labels)``; with ``supervision_verify`` + success also requires a fresh supervised PID ("the call returned" is not "supervised"). """ from hermes_cli.gateway import ( get_launchd_label, @@ -442,10 +433,10 @@ def _restart_launchd_gateway_after_update( if not supervision_verify: return [current_label], [] - # launchd_restart() returning only means "restart REQUESTED" (both branches are - # async). A helper dying before first bootstrap, or a bootstrap exiting 0 without - # registering (macOS 26.6.1), would otherwise reach "Update complete!" unsupervised. - # Verified domain-agnostically: domain locate fails on macOS-26 per-user domains. + # launchd_restart() returning only means "restart REQUESTED" (async). A helper dying + # before first bootstrap, or a bootstrap exiting 0 without registering (macOS 26.6.1), + # would otherwise reach "Update complete!" unsupervised. Verified domain-agnostically: + # domain locate fails on macOS-26 per-user domains. if wait_for_launchd_gateway_supervision(label=current_label): return [current_label], [] print( @@ -477,13 +468,11 @@ def _restart_macos_launchd_gateways( _wait_for_launchd_service_pid, ) - # --- Current profile --- _restarted, _failed = _restart_launchd_gateway_after_update(supervision_verify=True) restarted_services.extend(_restarted) failed_or_stale_units.extend(_failed) current_label = get_launchd_label() - # --- Sibling profiles --- for label in launchd_gateway_labels_for_install(): if label == current_label: continue @@ -492,8 +481,7 @@ def _restart_macos_launchd_gateways( # reuse that domain so a sibling is never probed in one and restarted in another. domain, old_pid = _locate_launchd_gateway_service(label) if domain is None: - # Installed but not bootstrapped — nothing running old code here. - continue + continue # installed but not bootstrapped — nothing running old code here graceful_ok = False if old_pid is not None and old_pid > 0: print(f" → {label}: draining (up to {int(drain_budget)}s)...") @@ -548,9 +536,6 @@ def _surviving_gateway_pids_after_failed_restart(): return None -_FRESH_RESTART_SUPERVISORS = frozenset({"systemd", "launchd", "service", "s6"}) - - def _gateway_service_matches_profile(profile: str, service: object) -> bool: """Match an exact gateway service/label (systemd/launchd/s6 shapes) to a profile. @@ -566,6 +551,23 @@ def _gateway_service_matches_profile(profile: str, service: object) -> bool: } +_MANUAL_GATEWAY_SKIP_REASON = ( + "manual gateway has no supervisor relaunch authority; left running for explicit operator restart" +) +_DESKTOP_SERVE_SKIP_REASON = ( + "desktop app owns and respawns this serve backend;" + " the recovery pass must not restart it out from under its supervisor" +) +# NOT a claim that no supervisor exists: a systemd-launched serve sets neither HERMES_SPAWN +# nor HERMES_PARENT_PID ("manual-serve"). Unit-backed serves are recovered by the fresh +# child's systemd pass; survivors reported by _surviving_pre_update_serve_runtimes. +_SERVE_SKIP_REASON = ( + "no per-profile relaunch command reaches a serve/dashboard runtime; recovered by the fresh" + " systemd unit pass when it owns a hermes-serve* unit, else left running for explicit" + " operator restart" +) + + def _gateway_recovery_partition( plan, *, skip_profiles: set[str] | None = None ) -> tuple[dict[str, str], list[dict]]: @@ -575,11 +577,9 @@ def _gateway_recovery_partition( failing interpreter is what raises the original ``ImportError``. Returns ``(candidates, skipped)``: profile → supervisor for supervised gateways the fresh process may restart; every other inventoried runtime with an explicit reason so - nothing vanishes silently (manual gateways have no relaunch authority; - serve/dashboard have no per-profile command). Skipped serve/dashboard is NOT - unrecoverable: the fresh child's ``hermes-serve*`` systemd pass enumerates units - from systemd (the ledger reads a systemd-launched serve as ``manual-serve``); - leftovers are caught by :func:`_surviving_pre_update_serve_runtimes`. + nothing vanishes silently. Skipped serve/dashboard is NOT unrecoverable: the fresh + child's ``hermes-serve*`` systemd pass enumerates units from systemd; leftovers are + caught by :func:`_surviving_pre_update_serve_runtimes`. """ skip_profiles = skip_profiles or set() candidates: dict[str, str] = {} @@ -596,45 +596,15 @@ def _gateway_recovery_partition( continue if supervisor in _FRESH_RESTART_SUPERVISORS: candidates.setdefault(profile, str(supervisor)) - else: - skipped.append( - { - "profile": profile, - "kind": "gateway", - "supervisor": str(supervisor), - "reason": ( - "manual gateway has no supervisor relaunch" - " authority; left running for explicit operator" - " restart" - ), - } - ) + continue + reason = _MANUAL_GATEWAY_SKIP_REASON elif kind in ("serve", "dashboard"): - if supervisor == "desktop": - reason = ( - "desktop app owns and respawns this serve backend;" - " the recovery pass must not restart it out from under" - " its supervisor" - ) - else: - # NOT a claim that no supervisor exists: a systemd-launched serve - # sets neither HERMES_SPAWN nor HERMES_PARENT_PID ("manual-serve"). - # Unit-backed serves are recovered by the fresh child's systemd pass; - # survivors reported by _surviving_pre_update_serve_runtimes. - reason = ( - "no per-profile relaunch command reaches a serve/" - "dashboard runtime; recovered by the fresh systemd" - " unit pass when it owns a hermes-serve* unit, else" - " left running for explicit operator restart" - ) - skipped.append( - { - "profile": profile, - "kind": str(kind), - "supervisor": str(supervisor), - "reason": reason, - } - ) + reason = _DESKTOP_SERVE_SKIP_REASON if supervisor == "desktop" else _SERVE_SKIP_REASON + else: + continue + skipped.append( + {"profile": profile, "kind": str(kind), "supervisor": str(supervisor), "reason": reason} + ) return candidates, skipped @@ -668,9 +638,8 @@ def _drain_or_signal_gateway_for_update( 1. Gateway is an ancestor of this process (auto-update cron inside the gateway tree): waiting is circular (gateway waits on in-flight work → cron session - waits on update → update waits on gateway); the wedged probe can't break it - (cron posts activity every ~180s) and the 1800s force-drain cap burns. So - fire-and-forget: signal restart and return; it completes once THIS process exits. + waits on update → update waits on gateway) and the 1800s force-drain cap burns. + So fire-and-forget: signal restart and return; it completes once THIS process exits. 2. Event loop provably wedged: SIGUSR1 can never drain it; bounded SIGTERM→SIGKILL. 3. Live out-of-tree gateway: graceful SIGUSR1 drain up to ``drain_budget``. """ @@ -714,24 +683,15 @@ def _resolve_manage_cmd(cache: dict, scope_: str, scope_cmd_: list, svc_name_: s if scope_ in cache: return cache[scope_] cmd = scope_cmd_ + ["--no-ask-password"] - if ( - scope_ == "system" - and hasattr(os, "geteuid") - and os.geteuid() != 0 # windows-footgun: ok — systemd path, Linux-only - ): - sudo_cmd = ["sudo", "-n"] + scope_cmd_ + ["--no-ask-password"] - sudo_ok = False + if _needs_sudo(scope_): + sudo_cmd = ["sudo", "-n"] + cmd try: - _probe = subprocess.run(["sudo", "-n", "true"], capture_output=True, timeout=5) - sudo_ok = _probe.returncode == 0 + sudo_ok = subprocess.run(["sudo", "-n", "true"], capture_output=True, timeout=5).returncode == 0 if not sudo_ok: # Blanket sudo refused — a targeted NOPASSWD sudoers entry may still work. - _probe = subprocess.run( - sudo_cmd + ["reset-failed", svc_name_], - capture_output=True, - timeout=5, - ) - sudo_ok = _probe.returncode == 0 + sudo_ok = subprocess.run( + sudo_cmd + ["reset-failed", svc_name_], capture_output=True, timeout=5 + ).returncode == 0 except (FileNotFoundError, subprocess.TimeoutExpired): sudo_ok = False cmd = sudo_cmd if sudo_ok else None @@ -753,7 +713,6 @@ def _restart_one_systemd_gateway_unit( Appends settled names to ``restarted_services`` and failures to ``failed_or_stale_units``. """ - # Check if active check = _systemctl(scope_cmd + ["is-active", svc_name], timeout=5) if check.stdout.strip() != "active": return @@ -762,60 +721,40 @@ def _restart_one_systemd_gateway_unit( # entirely or polkit prompts inside the captured subprocess. _manage_cmd = _resolve_manage_cmd(_manage_cmd_cache, scope, scope_cmd, svc_name) - # Graceful SIGUSR1 first so in-flight runs drain: handler → - # request_restart(via_service=True) → drain → exit, Restart=always - # respawns. hermes-serve has no handler → blunt restart below. + # Graceful SIGUSR1 first so in-flight runs drain: handler → request_restart(via_service=True) + # → drain → exit, Restart=always respawns. hermes-serve has no handler → blunt restart below. _main_pid = 0 if _service_unit_supports_graceful_sigusr1_restart(svc_name): try: _show = _systemctl(scope_cmd + ["show", svc_name, "--property=MainPID", "--value"], timeout=5) _main_pid = int((_show.stdout or "").strip() or 0) - except ( - ValueError, - subprocess.TimeoutExpired, - FileNotFoundError, - ): + except (ValueError, subprocess.TimeoutExpired, FileNotFoundError): _main_pid = 0 - _graceful_ok = False - if _main_pid > 0: - # Three-way triage (ancestor / wedged / graceful drain). - _graceful_ok = _drain_or_signal_gateway_for_update( - _main_pid, drain_budget, svc_name - ) + # Three-way triage (ancestor / wedged / graceful drain). + _graceful_ok = _main_pid > 0 and _drain_or_signal_gateway_for_update(_main_pid, drain_budget, svc_name) if _graceful_ok: - # ``Restart=always`` respawns only after RestartSec (60s in our unit; - # dead time for a voluntary restart). ``reset-failed`` + ``start`` - # skips it (~1-3s); if RestartSec already elapsed, ``start`` is a - # no-op and we fall through to the poll. Needs manage-units + # ``Restart=always`` respawns only after RestartSec (60s in our unit; dead time for a + # voluntary restart). ``reset-failed`` + ``start`` skips it (~1-3s); if RestartSec already + # elapsed, ``start`` is a no-op and we fall through to the poll. Needs manage-units # privileges; without them auto-restart still fires after RestartSec. if _manage_cmd is not None: _systemctl(_manage_cmd + ["reset-failed", svc_name], timeout=10) _systemctl(_manage_cmd + ["start", svc_name], timeout=15) - # RestartSec bypassed — should be up in seconds. - if _wait_for_service_active( - scope_cmd, - svc_name, - timeout=10.0, - ): + if _wait_for_service_active(scope_cmd, svc_name, timeout=10.0): restarted_services.append(svc_name) return # Passive poll: auto-restart fires after RestartSec regardless of # privileges — primary when _manage_cmd is None, fallback otherwise. _restart_sec = _service_restart_sec(scope_cmd, svc_name, default=0.0) - _post_drain_timeout = max(10.0, _restart_sec + 10.0) if _manage_cmd is None and _restart_sec > 5.0: print( f" → {svc_name}: waiting for systemd " f"auto-restart (~{int(_restart_sec)}s; " "no root for an immediate restart)..." ) - if _wait_for_service_active( - scope_cmd, - svc_name, - timeout=_post_drain_timeout, - ): + if _wait_for_service_active(scope_cmd, svc_name, timeout=max(10.0, _restart_sec + 10.0)): restarted_services.append(svc_name) return # Exited but not respawned (older unit without Restart=on-failure / @@ -835,44 +774,35 @@ def _restart_one_systemd_gateway_unit( ) return - # Blunt restart — only when the graceful path failed (no SIGUSR1 - # wiring, drain over budget, restart-policy mismatch). Mirrors - # `hermes gateway restart` (`systemd_restart()`). + # Blunt restart — only when the graceful path failed (no SIGUSR1 wiring, drain over + # budget, restart-policy mismatch). Mirrors `hermes gateway restart` (`systemd_restart()`). restart = _systemctl_reset_and_restart(_manage_cmd, svc_name) - if restart.returncode == 0: - # restart returns 0 even if the new process crashes at once — verify. - if _wait_for_service_active( - scope_cmd, - svc_name, - timeout=10.0, - ): - restarted_services.append(svc_name) - else: - # Retry once — transient startup failures (stale module cache, - # import race) often clear; reset-failed so the retry isn't blocked. - print(f" ⚠ {svc_name} died after restart, retrying...") - _systemctl_reset_and_restart(_manage_cmd, svc_name) - if _wait_for_service_active( - scope_cmd, - svc_name, - timeout=10.0, - ): - restarted_services.append(svc_name) - print(f" ✓ {svc_name} recovered on retry") - else: - failed_or_stale_units.append(svc_name) - _scope_flag = "--user " if scope == "user" else "" - _sudo_hint = "sudo " if scope == "system" else "" - print( - f" ✗ {svc_name} failed to stay running after restart.\n" - f" Check logs: {_sudo_hint}journalctl {_scope_flag}-u {svc_name} --since '2 min ago'\n" - f" Recover manually:\n" - f" {_sudo_hint}systemctl {_scope_flag}reset-failed {svc_name}\n" - f" {_sudo_hint}systemctl {_scope_flag}restart {svc_name}" - ) - else: + if restart.returncode != 0: failed_or_stale_units.append(svc_name) print(f" ⚠ Failed to restart {svc_name}: {restart.stderr.strip()}") + return + # restart returns 0 even if the new process crashes at once — verify. + if _wait_for_service_active(scope_cmd, svc_name, timeout=10.0): + restarted_services.append(svc_name) + return + # Retry once — transient startup failures (stale module cache, + # import race) often clear; reset-failed so the retry isn't blocked. + print(f" ⚠ {svc_name} died after restart, retrying...") + _systemctl_reset_and_restart(_manage_cmd, svc_name) + if _wait_for_service_active(scope_cmd, svc_name, timeout=10.0): + restarted_services.append(svc_name) + print(f" ✓ {svc_name} recovered on retry") + return + failed_or_stale_units.append(svc_name) + _scope_flag = "--user " if scope == "user" else "" + _sudo_hint = "sudo " if scope == "system" else "" + print( + f" ✗ {svc_name} failed to stay running after restart.\n" + f" Check logs: {_sudo_hint}journalctl {_scope_flag}-u {svc_name} --since '2 min ago'\n" + f" Recover manually:\n" + f" {_sudo_hint}systemctl {_scope_flag}reset-failed {svc_name}\n" + f" {_sudo_hint}systemctl {_scope_flag}restart {svc_name}" + ) def _restart_systemd_gateway_units( @@ -885,66 +815,52 @@ def _restart_systemd_gateway_units( """ from hermes_cli.gateway import supports_systemd_services, _ensure_user_systemd_env + if not supports_systemd_services(): + return _manage_cmd_cache: dict = {} + with suppress(Exception): + _ensure_user_systemd_env() - # hermes-gateway* (default + profiles) plus hermes-serve* (Desktop backend). - if supports_systemd_services(): - with suppress(Exception): - _ensure_user_systemd_env() + def _on_list_timeout(scope: str, exc: subprocess.TimeoutExpired) -> None: + # Discovery timeout — skip this scope, keep the other. + print( + f" ⚠ systemctl timed out listing {scope}-scope " + f"gateway units ({exc.cmd if exc.cmd else 'unknown command'}). " + f"Check the gateway with: hermes gateway status" + ) - for scope, scope_cmd in [ - ("user", ["systemctl", "--user"]), - ("system", ["systemctl"]), - ]: - try: - result = _systemctl( - scope_cmd + ["list-units", "hermes-gateway*", "hermes-serve*", - "--plain", "--no-legend", "--no-pager"], - timeout=10, - ) - except FileNotFoundError: - continue - except subprocess.TimeoutExpired as exc: - # Discovery timeout — skip this scope, keep the other. - print( - f" ⚠ systemctl timed out listing {scope}-scope " - f"gateway units ({exc.cmd if exc.cmd else 'unknown command'}). " - f"Check the gateway with: hermes gateway status" - ) - continue + def _on_unit_timeout(svc_name: str, exc: subprocess.TimeoutExpired) -> None: + # Isolate to this unit; a scope-wide handler used to abort every + # later gateway and leave the fleet on mixed code. + failed_or_stale_units.append(svc_name) + print( + f" ⚠ systemctl timed out restarting {svc_name} " + f"({exc.cmd if exc.cmd else 'unknown command'}); " + f"continuing with remaining gateways" + ) - def _on_unit_timeout(svc_name: str, exc: subprocess.TimeoutExpired) -> None: - # Isolate to this unit; a scope-wide handler used to abort every - # later gateway and leave the fleet on mixed code. - failed_or_stale_units.append(svc_name) - print( - f" ⚠ systemctl timed out restarting {svc_name} " - f"({exc.cmd if exc.cmd else 'unknown command'}); " - f"continuing with remaining gateways" - ) - - # Scope-qualify this scope's additions before the next scope can add a - # same-named unit; ``finally`` so a mid-scope abort keeps settled units. - _scope_mark = len(restarted_services) - try: - _for_each_systemd_gateway_unit( - result.stdout, - process_unit=lambda svc_name: _restart_one_systemd_gateway_unit( - svc_name, - scope=scope, - scope_cmd=scope_cmd, - drain_budget=drain_budget, - _manage_cmd_cache=_manage_cmd_cache, - restarted_services=restarted_services, - failed_or_stale_units=failed_or_stale_units, - ), - on_unit_timeout=_on_unit_timeout, - ) - finally: - restarted_scoped_units.update( - f"{scope}/{name}" - for name in restarted_services[_scope_mark:] - ) + for scope, scope_cmd, result in _systemd_gateway_unit_listings(_on_list_timeout): + # Scope-qualify this scope's additions before the next scope can add a + # same-named unit; ``finally`` so a mid-scope abort keeps settled units. + _scope_mark = len(restarted_services) + try: + _for_each_systemd_gateway_unit( + result.stdout, + process_unit=lambda svc_name: _restart_one_systemd_gateway_unit( + svc_name, + scope=scope, + scope_cmd=scope_cmd, + drain_budget=drain_budget, + _manage_cmd_cache=_manage_cmd_cache, + restarted_services=restarted_services, + failed_or_stale_units=failed_or_stale_units, + ), + on_unit_timeout=_on_unit_timeout, + ) + finally: + restarted_scoped_units.update( + f"{scope}/{name}" for name in restarted_services[_scope_mark:] + ) @dataclass @@ -961,25 +877,35 @@ class _GatewayRestartOutcome: externally_supervised_profiles: list killed_pids: set + def record_receipt(self, **extra) -> None: + """Best-effort ``record_gateway_restart`` from the current bookkeeping.""" + with suppress(Exception): + from hermes_cli.update_receipt import record_gateway_restart -def _restart_manual_gateways( - _drain_budget, - *, - restarted_services, - killed_pids, - relaunched_profiles, - externally_supervised_profiles, - find_gateway_pids, - find_profile_gateway_processes, - _prepare_profile_gateway_update_restart, - _get_service_pids, - _wait_for_gateway_exit, -): + record_gateway_restart( + restarted_services=self.restarted_services, + relaunched_profiles=self.relaunched_profiles, + externally_supervised_profiles=self.externally_supervised_profiles, + killed_pids=sorted(self.killed_pids), + failed_units=self.failed_or_stale_units, + incomplete=self.incomplete, + **extra, + ) + + +def _restart_manual_gateways(out: _GatewayRestartOutcome, _drain_budget) -> None: """Drain/stop every manual (non-service) gateway and print the restart summary. - Mutates the list/set args in place; raises so the caller's abort recovery fires. + Mutates ``out`` in place; raises so the caller's abort recovery fires. """ import signal as _signal + from hermes_cli.gateway import ( + find_gateway_pids, + find_profile_gateway_processes, + _prepare_profile_gateway_update_restart, + _get_service_pids, + _wait_for_gateway_exit, + ) # Exclude just-restarted service PIDs so we don't kill what systemd/launchd spawned. service_pids = _get_service_pids(all_profiles=True) @@ -1003,44 +929,40 @@ def _restart_manual_gateways( ) unrestartable_pids.add(pid) continue - # SIGUSR1 drain first, SIGTERM fallback if unsupported/over budget — the - # watcher relaunches either way. The helper announces its choice first - # because a silent full-budget wait reads as a hung update. - drained = _drain_or_signal_gateway_for_update(pid, _drain_budget, proc.profile) - if not drained: + # SIGUSR1 drain first, SIGTERM fallback if unsupported/over budget — the watcher + # relaunches either way. The helper announces its choice first because a silent + # full-budget wait reads as a hung update. + if not _drain_or_signal_gateway_for_update(pid, _drain_budget, proc.profile): with suppress(ProcessLookupError, PermissionError): os.kill(pid, _signal.SIGTERM) - # Wait ≤5s for exit: Telegram keeps the old getUpdates session ~30s; a new - # gateway inside that window gets a 409 (_handle_polling_conflict retries, - # but a brief wait avoids it on fast machines). + # Wait ≤5s for exit: Telegram keeps the old getUpdates session ~30s; a new gateway + # inside that window gets a 409 (_handle_polling_conflict retries, but a brief + # wait avoids it on fast machines). _wait_for_gateway_exit(timeout=5.0, force_after=None) - killed_pids.add(pid) + out.killed_pids.add(pid) if restart_mode == "external-supervisor": - externally_supervised_profiles.append(proc.profile) + out.externally_supervised_profiles.append(proc.profile) else: - relaunched_profiles.append(proc.profile) + out.relaunched_profiles.append(proc.profile) for pid in manual_pids: if pid in profile_processes and pid not in unrestartable_pids: continue with suppress(ProcessLookupError, PermissionError): os.kill(pid, _signal.SIGTERM) - killed_pids.add(pid) + out.killed_pids.add(pid) - if restarted_services or killed_pids: + if out.restarted_services or out.killed_pids: print() - for svc in restarted_services: + for svc in out.restarted_services: print(f" ✓ Restarted {svc}") - if relaunched_profiles: - names = ", ".join(relaunched_profiles) - print(f" ✓ Restarting manual gateway profile(s): {names}") - if externally_supervised_profiles: - names = ", ".join(externally_supervised_profiles) + if out.relaunched_profiles: + print(f" ✓ Restarting manual gateway profile(s): {', '.join(out.relaunched_profiles)}") + if out.externally_supervised_profiles: + names = ", ".join(out.externally_supervised_profiles) print(f" ✓ Handed gateway profile(s) back to their external supervisor: {names}") unmapped_count = ( - len(killed_pids) - - len(relaunched_profiles) - - len(externally_supervised_profiles) + len(out.killed_pids) - len(out.relaunched_profiles) - len(out.externally_supervised_profiles) ) if unmapped_count: print(f" → Stopped {unmapped_count} manual gateway process(es)") @@ -1049,54 +971,34 @@ def _restart_manual_gateways( print(" (or: hermes -p gateway run for each profile)") -def _force_kill_stuck_gateways(killed_pids, *, find_gateway_pids, _get_service_pids): - # Survivor sweep: gateways ignoring SIGTERM (stuck drain, blocked I/O, zombie) - # never exit, so the watcher never respawns and ImportErrors persist. Give - # graceful paths a moment, then SIGKILL remaining pre-update PIDs. +def _force_kill_stuck_gateways(killed_pids) -> None: + """Survivor sweep: gateways ignoring SIGTERM (stuck drain, blocked I/O, zombie) never + exit, so the watcher never respawns and ImportErrors persist. Give graceful paths a + moment, then SIGKILL remaining pre-update PIDs.""" with _best_effort('Post-restart survivor sweep failed: %s'): + from hermes_cli.gateway import find_gateway_pids, _get_service_pids + _time.sleep(3.0) - _service_pids_after = _get_service_pids(all_profiles=True) - _surviving = find_gateway_pids(exclude_pids=_service_pids_after, all_profiles=True) + _surviving = find_gateway_pids(exclude_pids=_get_service_pids(all_profiles=True), all_profiles=True) # Only PIDs we already tried to kill; newer ones are left alone. _stuck = [pid for pid in _surviving if pid in killed_pids] if _stuck: print() print(f" ⚠ {len(_stuck)} gateway process(es) ignored SIGTERM — force-killing") - from gateway.status import ( - get_process_start_time as _get_process_start_time, - terminate_pid as _terminate_pid, - ) + from gateway.status import get_process_start_time, terminate_pid + for pid in _stuck: with suppress(ProcessLookupError, PermissionError, OSError): # taskkill /T /F on Windows (no SIGKILL there), SIGKILL on POSIX. - _terminate_pid( - pid, - force=True, - expected_start_time=_get_process_start_time(pid), - ) + terminate_pid(pid, force=True, expected_start_time=get_process_start_time(pid)) # Let the OS reap so watchers see the exit and respawn. _time.sleep(1.5) def _recover_after_restart_phase_abort( - e, - _pre_update_plan, - *, - gateway_mode, - gateway_fleet_restart_incomplete, - gateway_restart_phase_errors, - _pre_restart_gateway_pids, - restarted_services, - restarted_scoped_units, - failed_or_stale_units, - relaunched_profiles, - externally_supervised_profiles, - killed_pids, -) -> bool: - """Phase-abort recovery: fresh-child restart + fail-closed verdict. - - Returns the new ``incomplete`` flag; extends the list args in place. - """ + e, _pre_update_plan, out: _GatewayRestartOutcome, *, gateway_mode, restarted_scoped_units +) -> None: + """Phase-abort recovery: fresh-child restart + fail-closed verdict; updates ``out`` in place.""" from hermes_cli.update_cmd import ( _abort_recovery_is_complete, _recover_gateway_restart_after_abort, @@ -1106,92 +1008,63 @@ def _recover_after_restart_phase_abort( ) logger.debug("Gateway restart during update failed: %s", e) - gateway_restart_phase_errors.append(str(e)) + out.phase_errors.append(str(e)) # Restart output never printed: assume stale unless provably no gateway runs. # Empty ``_surviving`` proves safety only if nothing ran beforehand; a gone # pre-restart gateway was stopped without verified replacement → fail closed. _surviving = _surviving_gateway_pids_after_failed_restart() - _already_restarted_profiles = set(relaunched_profiles) - _already_restarted_profiles.update(externally_supervised_profiles) - for runtime in getattr(_pre_update_plan, "runtimes", ()) or (): - if getattr(runtime, "kind", None) != "gateway": - continue - profile = getattr(runtime, "profile", None) - if not isinstance(profile, str): - continue - if any( - _gateway_service_matches_profile(profile, service) - for service in restarted_services - ): - _already_restarted_profiles.add(profile) + _planned_gateway_profiles = { + runtime.profile + for runtime in getattr(_pre_update_plan, "runtimes", ()) or () + if getattr(runtime, "kind", None) == "gateway" + and isinstance(getattr(runtime, "profile", None), str) + } + _already_restarted_profiles = set(out.relaunched_profiles) | set(out.externally_supervised_profiles) + _already_restarted_profiles.update( + profile + for profile in _planned_gateway_profiles + if any(_gateway_service_matches_profile(profile, service) for service in out.restarted_services) + ) _recovery_result = _recover_gateway_restart_after_abort( _pre_update_plan, gateway_mode=gateway_mode, skip_profiles=_already_restarted_profiles, skip_units=set(restarted_scoped_units), ) - _recovery_serve_units = _recovery_result.get("serve_units") or {} - _serve_units_failed = list(_recovery_serve_units.get("failed") or []) - # Deliberately NOT merged into ``restarted_services`` (gateway vocabulary feeding - # the fleet probe); serve coverage lives in the recovery result/receipt. A - # serve/dashboard still the SAME pre-update process is live on old code - # (unreachable by `gateway restart`): recovery may not claim success while one - # remains, and must never kill one (manual/Desktop serves have no relaunch authority). + _serve_units_failed = list((_recovery_result.get("serve_units") or {}).get("failed") or []) + # Deliberately NOT merged into ``restarted_services`` (gateway vocabulary feeding the + # fleet probe); serve coverage lives in the recovery result/receipt. A serve/dashboard + # still the SAME pre-update process is live on old code (unreachable by `gateway + # restart`): recovery may not claim success while one remains, and must never kill + # one (manual/Desktop serves have no relaunch authority). _stale_runtime_rows = _surviving_pre_update_serve_runtimes(_pre_update_plan) _recovery_result["stale_runtimes"] = _stale_runtime_rows # Only systemd-VERIFIED outcomes claim coverage; a relaunch that merely exited 0 # ("relaunch_attempted") was never observed and must not clear incomplete. _recovery_verified = set(_recovery_result.get("verified") or []) - if _recovery_verified: - relaunched_profiles.extend( - profile - for profile in sorted(_recovery_verified) - if profile not in relaunched_profiles - ) - _planned_gateway_runtimes = [ - runtime - for runtime in getattr(_pre_update_plan, "runtimes", ()) or () - if getattr(runtime, "kind", None) == "gateway" - and isinstance(getattr(runtime, "profile", None), str) - ] - _planned_gateway_profiles = {runtime.profile for runtime in _planned_gateway_runtimes} - _covered_gateway_profiles = _already_restarted_profiles | _recovery_verified - _recovery_complete = _abort_recovery_is_complete( + out.relaunched_profiles.extend( + profile for profile in sorted(_recovery_verified) if profile not in out.relaunched_profiles + ) + if _abort_recovery_is_complete( planned_gateway_profiles=_planned_gateway_profiles, - covered_gateway_profiles=_covered_gateway_profiles, + covered_gateway_profiles=_already_restarted_profiles | _recovery_verified, recovery_result=_recovery_result, stale_runtime_rows=_stale_runtime_rows, - ) - if _recovery_complete: + ): # Fresh child is terminal; the fleet-version matrix stays the authoritative # read-back before success is declared. - gateway_fleet_restart_incomplete = False + out.incomplete = False elif ( - _restart_phase_failure_is_incomplete( - _surviving, _pre_restart_gateway_pids - ) + _restart_phase_failure_is_incomplete(_surviving, out.pre_restart_gateway_pids) or _stale_runtime_rows or _serve_units_failed ): - gateway_fleet_restart_incomplete = True + out.incomplete = True _warn_gateway_restart_phase_aborted(e, _surviving) _warn_stale_serve_runtimes(_stale_runtime_rows) if gateway_mode: _write_gateway_update_exit_code(False) - with suppress(Exception): - from hermes_cli.update_receipt import record_gateway_restart - - record_gateway_restart( - restarted_services=restarted_services, - relaunched_profiles=relaunched_profiles, - externally_supervised_profiles=externally_supervised_profiles, - killed_pids=sorted(killed_pids), - failed_units=failed_or_stale_units, - incomplete=gateway_fleet_restart_incomplete, - phase_error=str(e), - fresh_recovery=_recovery_result, - ) - return gateway_fleet_restart_incomplete + out.record_receipt(phase_error=str(e), fresh_recovery=_recovery_result) def _restart_gateway_fleet_after_update(_pre_update_plan, gateway_mode: bool): @@ -1202,30 +1075,33 @@ def _restart_gateway_fleet_after_update(_pre_update_plan, gateway_mode: bool): """ from hermes_cli.update_cmd import _m, _write_gateway_update_exit_code - gateway_fleet_restart_incomplete = False - gateway_restart_phase_errors: list[str] = [] - # Empty until we are about to stop/drain, so an early exception has nothing to - # fail closed on, while a failure after stopping a discovered gateway fails - # closed on an empty survivor probe. - _pre_restart_gateway_pids: list | None = [] - # All bookkeeping is declared before the try (never reset to None) so abort - # recovery and fleet reconciliation can read it even if the phase raises early. - restarted_services: list = [] - # Scope-qualified twin (``user/hermes-serve`` vs ``system/hermes-serve`` are - # different processes; abort recovery needs WHICH settled). Bare names stay in + # All bookkeeping is declared before the try so abort recovery and fleet reconciliation + # can read it even if the phase raises early. ``pre_restart_gateway_pids`` stays empty + # until we are about to stop/drain, so an early exception has nothing to fail closed on, + # while a failure after stopping a discovered gateway fails closed on an empty survivor probe. + out = _GatewayRestartOutcome( + incomplete=False, + phase_errors=[], + pre_restart_gateway_pids=[], + restarted_services=[], + failed_or_stale_units=[], + relaunched_profiles=[], + externally_supervised_profiles=[], + killed_pids=set(), + ) + # Scope-qualified twin (``user/hermes-serve`` vs ``system/hermes-serve`` are different + # processes; abort recovery needs WHICH settled). Bare names stay in # ``restarted_services`` for the fleet probe, receipt and summary. restarted_scoped_units: set = set() - failed_or_stale_units: list = [] - relaunched_profiles: list = [] - externally_supervised_profiles: list = [] - killed_pids: set = set() # Purge stale cached Hermes modules FIRST: the import below loads new gateway # source into this pre-update interpreter, and a cached sibling missing a # symbol the new source expects would ImportError and abort the whole phase. _m()._purge_stale_hermes_modules() try: - from hermes_cli.gateway import ( + # Every gateway helper the phase needs is imported up front so a broken gateway + # module aborts into recovery BEFORE any unit is touched. + from hermes_cli.gateway import ( # noqa: F401 is_macos, find_gateway_pids, find_profile_gateway_processes, @@ -1244,93 +1120,90 @@ def _restart_gateway_fleet_after_update(_pre_update_plan, gateway_mode: bool): except Exception: _drain_budget = 45.0 - failed_or_stale_units = [] - killed_pids = set() - relaunched_profiles = [] - externally_supervised_profiles = [] - # Snapshot before any stop/drain so an empty survivor probe reads as "stopped # and never came back", not "nothing was running"; None fails closed. try: - _pre_restart_gateway_pids = list(find_gateway_pids(all_profiles=True)) + out.pre_restart_gateway_pids = list(find_gateway_pids(all_profiles=True)) except Exception: - _pre_restart_gateway_pids = None + out.pre_restart_gateway_pids = None _restart_systemd_gateway_units( - restarted_services, failed_or_stale_units, restarted_scoped_units, _drain_budget + out.restarted_services, out.failed_or_stale_units, restarted_scoped_units, _drain_budget ) # macOS: EVERY ai.hermes.gateway* LaunchAgent (systemd parity). if is_macos(): with suppress(FileNotFoundError, ImportError): _restart_macos_launchd_gateways( - restarted_services, - failed_or_stale_units, - _drain_budget, + out.restarted_services, out.failed_or_stale_units, _drain_budget ) - _restart_manual_gateways( - _drain_budget, - restarted_services=restarted_services, - killed_pids=killed_pids, - relaunched_profiles=relaunched_profiles, - externally_supervised_profiles=externally_supervised_profiles, - find_gateway_pids=find_gateway_pids, - find_profile_gateway_processes=find_profile_gateway_processes, - _prepare_profile_gateway_update_restart=_prepare_profile_gateway_update_restart, - _get_service_pids=_get_service_pids, - _wait_for_gateway_exit=_wait_for_gateway_exit, - ) + _restart_manual_gateways(out, _drain_budget) - if failed_or_stale_units: - gateway_fleet_restart_incomplete = True + if out.failed_or_stale_units: + out.incomplete = True if gateway_mode: _write_gateway_update_exit_code(False) - _warn_incomplete_gateway_fleet_restart(failed_or_stale_units) - - with suppress(Exception): - from hermes_cli.update_receipt import record_gateway_restart - - record_gateway_restart( - restarted_services=restarted_services, - relaunched_profiles=relaunched_profiles, - externally_supervised_profiles=externally_supervised_profiles, - killed_pids=sorted(killed_pids), - failed_units=failed_or_stale_units, - incomplete=bool(failed_or_stale_units), - ) - - _force_kill_stuck_gateways( - killed_pids, find_gateway_pids=find_gateway_pids, _get_service_pids=_get_service_pids - ) + _warn_incomplete_gateway_fleet_restart(out.failed_or_stale_units) + out.record_receipt() + _force_kill_stuck_gateways(out.killed_pids) except Exception as e: - gateway_fleet_restart_incomplete = _recover_after_restart_phase_abort( - e, - _pre_update_plan, - gateway_mode=gateway_mode, - gateway_fleet_restart_incomplete=gateway_fleet_restart_incomplete, - gateway_restart_phase_errors=gateway_restart_phase_errors, - _pre_restart_gateway_pids=_pre_restart_gateway_pids, - restarted_services=restarted_services, - restarted_scoped_units=restarted_scoped_units, - failed_or_stale_units=failed_or_stale_units, - relaunched_profiles=relaunched_profiles, - externally_supervised_profiles=externally_supervised_profiles, - killed_pids=killed_pids, + _recover_after_restart_phase_abort( + e, _pre_update_plan, out, gateway_mode=gateway_mode, restarted_scoped_units=restarted_scoped_units ) - return _GatewayRestartOutcome( - incomplete=gateway_fleet_restart_incomplete, - phase_errors=gateway_restart_phase_errors, - pre_restart_gateway_pids=_pre_restart_gateway_pids, - restarted_services=restarted_services, - failed_or_stale_units=failed_or_stale_units, - relaunched_profiles=relaunched_profiles, - externally_supervised_profiles=externally_supervised_profiles, - killed_pids=killed_pids, + return out + + +def _print_legacy_units_warning() -> None: + """Legacy hermes.service fights hermes-gateway.service over the bot token; warn on + every update until migrated.""" + from hermes_cli.gateway import ( + has_legacy_hermes_units, + _find_legacy_hermes_units, + supports_systemd_services, ) + if not (supports_systemd_services() and has_legacy_hermes_units()): + return + print() + print("⚠ Legacy Hermes gateway unit(s) detected:") + for name, path, is_sys in _find_legacy_hermes_units(): + scope = "system" if is_sys else "user" + print(f" {path} ({scope} scope)") + print() + print(" These pre-rename units (hermes.service) fight the current") + print(" hermes-gateway.service for the bot token and cause SIGTERM") + print(" flap loops. Remove them with:") + print() + print(" hermes gateway migrate-legacy") + print() + print(" (add `sudo` if any are in system scope)") + + +def _collect_fleet_snapshot(restart, rows_expected: bool) -> list: + """Fleet version rows, polled over a bounded settle window when runtimes are expected. + + Gateways need time to rewrite gateway_state.json; Windows resumes DETACHED (~10s boot), + so a single 2s sleep reported "no rows" on healthy resumes. A "down" row may be a + detached replacement still booting: poll until none remain or the deadline passes. + Pre-restart PIDs make a gateway stopped WITHOUT verified replacement a DOWN row (exit 1) + instead of no row at all. + """ + from hermes_cli.update_receipt import collect_fleet_versions + + if not rows_expected: + return collect_fleet_versions(pre_restart_pids=restart.pre_restart_gateway_pids) + _fleet_deadline = _time.monotonic() + 30.0 + while True: + _time.sleep(2.0) + snapshot = collect_fleet_versions(pre_restart_pids=restart.pre_restart_gateway_pids) + if snapshot and not any(row.get("state") == "down" for row in snapshot): + return snapshot + if _time.monotonic() >= _fleet_deadline: + return snapshot + def _verify_fleet_after_update( restart, @@ -1352,29 +1225,8 @@ def _verify_fleet_after_update( _surviving_pre_update_serve_runtimes, _warn_stale_serve_runtimes, ) - # Legacy hermes.service fights hermes-gateway.service over the bot token; warn - # on every update until migrated. with _best_effort('Legacy unit check during update failed: %s'): - from hermes_cli.gateway import ( - has_legacy_hermes_units, - _find_legacy_hermes_units, - supports_systemd_services, - ) - - if supports_systemd_services() and has_legacy_hermes_units(): - print() - print("⚠ Legacy Hermes gateway unit(s) detected:") - for name, path, is_sys in _find_legacy_hermes_units(): - scope = "system" if is_sys else "user" - print(f" {path} ({scope} scope)") - print() - print(" These pre-rename units (hermes.service) fight the current") - print(" hermes-gateway.service for the bot token and cause SIGTERM") - print(" flap loops. Remove them with:") - print() - print(" hermes gateway migrate-legacy") - print() - print(" (add `sudo` if any are in system scope)") + _print_legacy_units_warning() # Restart a managed dashboard via systemd or stop stale manual ones (raw-killing # a systemd-owned PID reads as clean stop and leaves the Cloudflare origin dead). @@ -1401,10 +1253,7 @@ def _verify_fleet_after_update( # instead of assuming the restart phase worked. _fleet_snapshot: list = [] with _best_effort('Fleet version verification failed: %s'): - from hermes_cli.update_receipt import ( - collect_fleet_versions, - print_fleet_version_matrix, - ) + from hermes_cli.update_receipt import print_fleet_version_matrix # Cross-platform "rows expected" signal: (restarted_services or killed_pids) # never fires on Windows (pause/resume populates neither), so a healthy @@ -1416,31 +1265,7 @@ def _verify_fleet_after_update( restart.restarted_services, restart.killed_pids, ) - # Settle window (skipped when nothing ran): gateways need time to rewrite - # gateway_state.json. Windows resumes DETACHED (~10s boot), so a single 2s - # sleep reported "no rows" on healthy resumes; poll a bounded window instead. - _fleet_snapshot = [] - if _fleet_rows_expected: - _fleet_deadline = _time.monotonic() + 30.0 - while True: - _time.sleep(2.0) - # Pre-restart PIDs make a gateway stopped WITHOUT verified - # replacement a DOWN row (exit 1) instead of no row at all. - _fleet_snapshot = collect_fleet_versions( - pre_restart_pids=restart.pre_restart_gateway_pids - ) - # A "down" row may be a detached replacement still booting; - # poll until no "down" rows remain or the deadline passes. - if _fleet_snapshot and not any( - row.get("state") == "down" for row in _fleet_snapshot - ): - break - if _time.monotonic() >= _fleet_deadline: - break - else: - _fleet_snapshot = collect_fleet_versions( - pre_restart_pids=restart.pre_restart_gateway_pids - ) + _fleet_snapshot = _collect_fleet_snapshot(restart, _fleet_rows_expected) if print_fleet_version_matrix(_fleet_snapshot): restart.incomplete = True elif not _fleet_snapshot and _fleet_rows_expected: @@ -1454,7 +1279,6 @@ def _verify_fleet_after_update( # Every runtime the PLAN saw must appear in restart bookkeeping; an # unaccounted one is a silent miss and escalates like a STALE/DOWN row. - _runtime_outcomes: list = [] with _best_effort('Runtime-outcome reconciliation failed: %s'): if _pre_update_plan is not None and _pre_update_plan.runtimes: from hermes_cli.update_inventory import ( @@ -1488,11 +1312,7 @@ def _verify_fleet_after_update( from hermes_cli.update_receipt import finalize_update_receipt _receipt_path = finalize_update_receipt( - ( - "partial" - if restart.incomplete or not update_complete - else "success" - ), + "partial" if restart.incomplete or not update_complete else "success", fleet=_fleet_snapshot, ) if _receipt_path is not None: @@ -1526,20 +1346,16 @@ def _fleet_probe_expected_runtimes( ) -> bool: """Whether the post-update fleet probe should have produced rows. - ``collect_fleet_versions()`` swallows every failure and an empty matrix prints - as healthy, so zero rows is only proof-of-safety when NOTHING says a gateway - existed pre-update. Signals: ``restarted_services``/``killed_pids`` (POSIX phase - touched gateways); ``pre_restart_pids`` non-empty or None (unreadable, same - contract as ``_restart_phase_failure_is_incomplete``); plan inventoried ≥1 runtime. - - ``windows_resume_token`` is deliberately EXCLUDED: it is pause/resume bookkeeping, - not an inventory, and its entries don't map to probe rows (``unmapped`` - Scheduled-Task gateways never publish gateway_state.json; a paused profile resumes - DETACHED and may not republish within the window). Counting it made every - Windows update that paused a gateway exit 1 after a long silent wait. A live - pre-update Windows gateway is already covered by ``pre_restart_pids`` and the plan. - The same condition gates the settle sleep. Non-empty snapshots (incl. ``unknown`` - rows) are judged solely by ``print_fleet_version_matrix``. + ``collect_fleet_versions()`` swallows every failure and an empty matrix prints as + healthy, so zero rows is only proof-of-safety when NOTHING says a gateway existed + pre-update. Signals: ``restarted_services``/``killed_pids``; ``pre_restart_pids`` + non-empty or None (same contract as ``_restart_phase_failure_is_incomplete``); plan + inventoried ≥1 runtime. ``windows_resume_token`` is deliberately EXCLUDED: it is + pause/resume bookkeeping, not an inventory, and its entries don't map to probe rows + (``unmapped`` Scheduled-Task gateways never publish gateway_state.json; a paused + profile resumes DETACHED). Counting it made every Windows update that paused a + gateway exit 1 after a long silent wait; a live pre-update Windows gateway is already + covered by ``pre_restart_pids`` and the plan. The same condition gates the settle sleep. """ del windows_resume_token # excluded on purpose — see docstring if restarted_services or killed_pids: @@ -1570,6 +1386,9 @@ def _wait_for_service_active( _time.sleep(0.5) +_RESTART_SEC_UNITS = (("ms", 0.001), ("us", 0.000001), ("min", 60.0), ("s", 1.0)) + + def _service_restart_sec( scope_cmd_: list, svc_name_: str, @@ -1588,12 +1407,7 @@ def _service_restart_sec( total = 0.0 matched = False for part in raw.split(): - for _suf, _mult in ( - ("ms", 0.001), - ("us", 0.000001), - ("min", 60.0), - ("s", 1.0), - ): + for _suf, _mult in _RESTART_SEC_UNITS: if part.endswith(_suf): with suppress(ValueError): total += float(part[: -len(_suf)]) * _mult diff --git a/hermes_cli/update_cmd_maint.py b/hermes_cli/update_cmd_maint.py index 30c44c8b2a..9de3ffb06f 100644 --- a/hermes_cli/update_cmd_maint.py +++ b/hermes_cli/update_cmd_maint.py @@ -5,6 +5,7 @@ still resolves/monkeypatches. Origin helpers are imported lazily per function (n test patches on ``update_cmd`` stay effective). """ +import importlib import logging from contextlib import suppress import os @@ -24,12 +25,10 @@ logger = logging.getLogger("hermes_cli.update_cmd") _UPDATE_RUNTIME_RELOAD_MODULES = "hermes_constants", "tools.environments.local", "tools.lazy_deps" - #: Package prefixes whose cached modules go stale when the checkout changes under this #: process; purged (not reloaded) so any LATER import chain resolves against fresh source. _STALE_PURGE_PREFIXES = "hermes_cli", "gateway", "tools", "tui_gateway", "agent" - #: Modules EXECUTING the update survive the purge: evicting them buys nothing (running frames #: keep them alive) and reloading them mid-flight is the one genuinely unsafe move. _STALE_PURGE_PROTECTED = frozenset( @@ -45,6 +44,36 @@ _STALE_PURGE_PROTECTED = frozenset( )} ) +_PRE_UPDATE_SNAPSHOT_KEEP = 1 + +# Per-file cap for the quick snapshot (larger files skipped with a warning): it protects +# small hard-to-regenerate state, not a multi-GB state.db (24 GB cost ~60s + 24 GB/update). +_PRE_UPDATE_SNAPSHOT_MAX_FILE_SIZE = 1 << 30 # 1 GiB + +_SQLITE_WAL_BUG_DETAIL = "SQLite {} still has the WAL-reset corruption bug" + + +def _load_updates_cfg() -> dict: + """``updates`` section of config.yaml; ``{}`` on any failure.""" + from hermes_cli.config import load_config + + cfg = load_config() or {} + updates = cfg.get("updates", {}) if isinstance(cfg, dict) else {} + return updates if isinstance(updates, dict) else {} + + +def _reload_modules(names, *, modules, log) -> None: + """``importlib.reload`` each module of *names* cached in *modules*; failures go to *log*.""" + importlib.invalidate_caches() + for module_name in names: + module = modules.get(module_name) + if module is None: + continue + try: + importlib.reload(module) + except Exception as exc: + log(module_name, exc) + def _purge_stale_hermes_modules() -> None: """Evict every cached Hermes module after the checkout changed in-place. Never raises. @@ -56,21 +85,15 @@ def _purge_stale_hermes_modules() -> None: """ from hermes_cli.update_cmd import _m with _best_effort('Could not purge stale Hermes modules: %s'): - import importlib - importlib.invalidate_caches() - purged = [] - for name in list(_m().sys.modules): - if name in _STALE_PURGE_PROTECTED: - continue - if not name.startswith(_STALE_PURGE_PREFIXES): - continue - root = name.split(".", 1)[0] - if root not in _STALE_PURGE_PREFIXES: - # startswith() also matches unrelated packages like ``gateway_foo``. - continue - if _m().sys.modules.pop(name, None) is not None: - purged.append(name) + modules = _m().sys.modules + purged = [ + name for name in list(modules) + if name not in _STALE_PURGE_PROTECTED + # Root-package check: startswith() alone also matches unrelated ``gateway_foo``. + and name.split(".", 1)[0] in _STALE_PURGE_PREFIXES + and modules.pop(name, None) is not None + ] if purged: logger.debug("Purged %d stale Hermes module(s) after checkout update", len(purged)) @@ -80,17 +103,11 @@ def _reload_updated_runtime_modules() -> None: can expose old symbols despite new source on disk.""" from hermes_cli.update_cmd import _m with _best_effort('Could not refresh update runtime modules: %s'): - import importlib - - importlib.invalidate_caches() - for module_name in _UPDATE_RUNTIME_RELOAD_MODULES: - module = _m().sys.modules.get(module_name) - if module is None: - continue - try: - importlib.reload(module) - except Exception as exc: - logger.debug("Could not reload updated module %s: %s", module_name, exc) + _reload_modules( + _UPDATE_RUNTIME_RELOAD_MODULES, + modules=_m().sys.modules, + log=lambda name, exc: logger.debug("Could not reload updated module %s: %s", name, exc), + ) def _print_curator_first_run_notice() -> None: @@ -99,9 +116,7 @@ def _print_curator_first_run_notice() -> None: disable it first. Silent on steady state.""" try: from agent import curator - except Exception: - return - try: + if not curator.is_enabled(): return state = curator.load_state() @@ -131,14 +146,11 @@ def _print_fts_optimize_available_notice() -> None: ``sessions.fts_optimize_notice``: ``advise`` (default), ``require`` (firmer), ``off``. """ - mode = "advise" try: from hermes_cli.config import load_config mode = str( - ((load_config() or {}).get("sessions") or {}).get( - "fts_optimize_notice", "advise" - ) + ((load_config() or {}).get("sessions") or {}).get("fts_optimize_notice", "advise") ).strip().lower() except Exception: mode = "advise" @@ -161,7 +173,6 @@ def _print_fts_optimize_available_notice() -> None: if size_gb < 0.5: return db = None - interrupted = False try: db = SessionDB(db_path=db_path, read_only=True) # read_only opens skip schema init; probe the layout directly. @@ -235,9 +246,7 @@ def _print_curator_recent_run_notice() -> None: ``last_run_summary_shown_at``. Silent when never run, already shown, or no rename info.""" try: from agent import curator - except Exception: - return - try: + state = curator.load_state() except Exception: return @@ -245,27 +254,19 @@ def _print_curator_recent_run_notice() -> None: last_run_at = state.get("last_run_at") if not last_run_at: return # no curator run yet — first-run notice handles this case - if state.get("last_run_summary_shown_at") == last_run_at: return # already shown for this run - summary = state.get("last_run_summary") or "" if not summary: return # Only a multi-line summary (rename map) is worth showing; still stamp it shown. - if "\n" not in summary: - with suppress(Exception): - state["last_run_summary_shown_at"] = last_run_at - curator.save_state(state) - return - - when = _format_time_ago(last_run_at) - print() - print(f"ℹ Skill curator — last run {when}") - for line in summary.splitlines(): - print(f" {line}") - print(" (This message shows once per curator run. View anytime: hermes curator status)") + if "\n" in summary: + print() + print(f"ℹ Skill curator — last run {_format_time_ago(last_run_at)}") + for line in summary.splitlines(): + print(f" {line}") + print(" (This message shows once per curator run. View anytime: hermes curator status)") with suppress(Exception): state["last_run_summary_shown_at"] = last_run_at @@ -279,8 +280,7 @@ def _format_time_ago(iso_ts: str) -> str: ts = datetime.fromisoformat(iso_ts.replace("Z", "+00:00")) if ts.tzinfo is None: ts = ts.replace(tzinfo=timezone.utc) - delta = datetime.now(timezone.utc) - ts - secs = int(delta.total_seconds()) + secs = int((datetime.now(timezone.utc) - ts).total_seconds()) if secs < 60: return "just now" if secs < 3600: @@ -297,20 +297,14 @@ def _reload_process_scan_modules() -> None: fresh ``_subprocess_compat``: cleanup runs in the PRE-update process and a symbol the update added would otherwise ImportError after the code update succeeded. Called from the cleanup entry point so every caller (git path, ZIP fallback) is covered.""" - import importlib - - importlib.invalidate_caches() - for mod_name in ( - "hermes_cli._subprocess_compat", - "hermes_cli.dashboard_procs", - ): - mod = sys.modules.get(mod_name) - if mod is not None: - try: - importlib.reload(mod) - except Exception as exc: - # warning, not debug: a failed reload surfaces as ImportError seconds later. - logger.warning("Could not reload %s for post-update cleanup: %s", mod_name, exc) + _reload_modules( + ("hermes_cli._subprocess_compat", "hermes_cli.dashboard_procs"), + modules=sys.modules, + # warning, not debug: a failed reload surfaces as ImportError seconds later. + log=lambda name, exc: logger.warning( + "Could not reload %s for post-update cleanup: %s", name, exc + ), + ) def _finish_dashboard_update_cleanup( @@ -407,17 +401,11 @@ def _print_verified_update_completion(message: str) -> bool: # Grace path: an unprobeable interpreter (dev checkout, no probe subprocess) must not # fail the update — only a POSITIVE vulnerable probe withholds success. logger.debug("Post-update SQLite runtime probe unavailable; not blocking") - _print_update_completion(message) - return True - if sqlite_runtime_ok: + if sqlite_info is None or sqlite_runtime_ok: _print_update_completion(message) return True print() - detail = ( - f"SQLite {sqlite_info.sqlite_version_string} still has the " - "WAL-reset corruption bug" - ) - print(f"⚠ Update partially complete — {detail}.") + print(f"⚠ Update partially complete — {_SQLITE_WAL_BUG_DETAIL.format(sqlite_info.sqlite_version_string)}.") print( " Rebuild the Hermes venv with a uv-managed Python, restart Hermes, " "then verify with `hermes doctor`." @@ -458,10 +446,7 @@ def _print_update_summary( if not desktop_build_ok: parts.append("the desktop app was not rebuilt and is still on the previous build") if not sqlite_runtime_ok and sqlite_info is not None: - parts.append( - f"SQLite {sqlite_info.sqlite_version_string} still has the " - "WAL-reset corruption bug" - ) + parts.append(_SQLITE_WAL_BUG_DETAIL.format(sqlite_info.sqlite_version_string)) print("⚠ Update partially complete — " + "; ".join(parts) + ".") if node_failures: print(" Code and Python deps are updated, but the dashboard/TUI may") @@ -539,14 +524,11 @@ def _verify_and_restore_one_state_db(home: Path, *, label: str) -> None: if not snap_root.exists(): print(" ⚠ No pre-update snapshot for this home") return - for snap_dir in sorted( - (d for d in snap_root.iterdir() if d.is_dir()), reverse=True - ): + for snap_dir in sorted((d for d in snap_root.iterdir() if d.is_dir()), reverse=True): snap_state = snap_dir / "state.db" if not snap_state.exists(): continue - snap_ok = verify_sqlite_integrity(snap_state, check_header=True, run_pragma=True) - if not snap_ok.get("valid"): + if not verify_sqlite_integrity(snap_state, check_header=True, run_pragma=True).get("valid"): continue try: if _restore_state_db_from_snapshot(state_path, snap_state): @@ -650,13 +632,10 @@ def _ensure_fhs_path_guard() -> None: except OSError: continue # Idempotency: any uncommented PATH line referencing /usr/local/bin (install.sh grep). - already_guarded = any( - "/usr/local/bin" in line - and "PATH" in line - and not line.lstrip().startswith("#") + if any( + "/usr/local/bin" in line and "PATH" in line and not line.lstrip().startswith("#") for line in existing.splitlines() - ) - if already_guarded: + ): continue try: with cfg.open("a", encoding="utf-8") as f: @@ -676,12 +655,11 @@ def _ensure_acp_launcher() -> None: delegates to the sibling ``hermes acp``, correct for every layout. No-op on Windows (install.ps1 stages launchers into ``$HermesHome\bin``, never - ``venv\Scripts`` which would shadow the user's python) and where it already exists. - Unwritable dirs are skipped. Idempotent. + ``venv\Scripts`` which would shadow the user's python; launcher repair lives in + _install_repair) and where it already exists. Unwritable dirs are skipped. Idempotent. """ from hermes_cli.update_cmd import _m if _m().sys.platform == "win32": - # Windows launcher repair lives in _install_repair, not here. return for bin_dir in (Path.home() / ".local" / "bin", Path("/usr/local/bin")): hermes_cmd = bin_dir / "hermes" @@ -706,12 +684,11 @@ def _ensure_acp_launcher() -> None: print(f" ✓ Installed hermes-acp launcher → {acp_cmd}") -_PRE_UPDATE_SNAPSHOT_KEEP = 1 - - -# Per-file cap for the quick snapshot (larger files skipped with a warning): it protects -# small hard-to-regenerate state, not a multi-GB state.db (24 GB cost ~60s + 24 GB/update). -_PRE_UPDATE_SNAPSHOT_MAX_FILE_SIZE = 1 << 30 # 1 GiB +_BACKUP_MODE_ALIASES = { + "off": "off", "false": "off", "none": "off", "disabled": "off", + "full": "full", "zip": "full", "true": "full", + "quick": "quick", +} def _resolve_pre_update_backup_mode(args) -> str: @@ -724,139 +701,102 @@ def _resolve_pre_update_backup_mode(args) -> str: return "full" try: - from hermes_cli.config import load_config - - cfg = load_config() + raw = _load_updates_cfg().get("pre_update_backup", "quick") except Exception as exc: logger.debug("Could not load config for pre-update backup: %s", exc) - cfg = {} - - updates_cfg = cfg.get("updates", {}) if isinstance(cfg, dict) else {} - raw = updates_cfg.get("pre_update_backup", "quick") + raw = "quick" if raw is True: return "full" if raw is False: return "off" - mode = str(raw).strip().lower() - if mode in ("off", "false", "none", "disabled"): - return "off" - if mode in ("full", "zip", "true"): - return "full" - if mode == "quick": + mode = _BACKUP_MODE_ALIASES.get(str(raw).strip().lower()) + if mode is None: + logger.warning("Unknown updates.pre_update_backup value %r — using 'quick'", raw) return "quick" - logger.warning("Unknown updates.pre_update_backup value %r — using 'quick'", raw) - return "quick" + return mode -def _run_pre_update_backup(args) -> Optional[str]: - """Run the pre-update backup; return the quick-snapshot id (None when off/failed). Never raises. +def _verify_state_db_after_snapshot(snapshot_id: str) -> None: + """Verify live state.db after the snapshot: a concurrent process (antivirus, killed + gateway, Windows filter driver) can corrupt it and we'd otherwise exit 0 silently.""" + from hermes_cli.backup import _quick_snapshot_root, verify_sqlite_integrity + from hermes_cli.config import get_hermes_home - ``off`` — nothing. ``quick`` (default) — snapshot of critical small files under - ``state-snapshots/``, files over 1 GiB skipped so a bloated state.db can't stall the update. - ``full`` — quick snapshot PLUS a zip of HERMES_HOME under ``backups/`` (``hermes import``). - """ - from hermes_cli.update_cmd import _record_update_step - mode = _resolve_pre_update_backup_mode(args) - - if mode == "off": - if getattr(args, "no_backup", False): - print("◆ Pre-update backup: skipped (--no-backup)") - print() - # Config-level off is silent: the user opted out. - return None - - snapshot_id = None - with _best_effort('Pre-update snapshot failed: %s'): - from hermes_cli.backup import ( - _quick_snapshot_root, - create_quick_snapshot, - verify_sqlite_integrity, + _src_path = get_hermes_home() / "state.db" + if not _src_path.exists(): + return + _integrity = verify_sqlite_integrity( + _src_path, + check_header=True, + run_pragma=True, + max_bytes=_PRE_UPDATE_SNAPSHOT_MAX_FILE_SIZE, + ) + if _integrity.get("valid"): + return + print(f" ⚠ state.db integrity check FAILED after snapshot: {_integrity.get('message', 'unknown error')}") + _snap_state = _quick_snapshot_root(get_hermes_home()) / snapshot_id / "state.db" + if not _snap_state.exists(): + print(" ⚠ Snapshot does not contain state.db (was skipped or too large).") + elif verify_sqlite_integrity(_snap_state, check_header=True, run_pragma=True).get("valid"): + print(" ✓ Snapshot copy is valid — continuing update.") + print(" If state.db is lost after update it will be auto-restored.") + else: + print( + " ✗ Snapshot copy ALSO failed integrity — " + "the source was already corrupted before the backup." ) + print() - # A later `from hermes_constants import get_hermes_home` in this function makes that - # name function-local (unbound here), so alias explicitly. - from hermes_cli.config import get_hermes_home as _get_home - snapshot_id = create_quick_snapshot( - label="pre-update", +def _run_quick_snapshots() -> Optional[str]: + """Quick snapshot of the root home plus every sibling profile; returns the root snapshot id.""" + from hermes_cli.update_cmd import _record_update_step + from hermes_cli.backup import create_quick_snapshot + + snapshot_id = create_quick_snapshot( + label="pre-update", + keep=_PRE_UPDATE_SNAPSHOT_KEEP, + max_file_size=_PRE_UPDATE_SNAPSHOT_MAX_FILE_SIZE, + ) + if snapshot_id: + _verify_state_db_after_snapshot(snapshot_id) + print(f"◆ Pre-update snapshot: {snapshot_id}") + + # The code swap + fleet restart touch EVERY profile, so each gets the same snapshot + # under its own state-snapshots/. Best-effort per profile. + with _best_effort('Sibling profile snapshots failed: %s'): + from hermes_cli.backup import create_pre_update_snapshots_all_profiles + + _sibling_snaps = create_pre_update_snapshots_all_profiles( keep=_PRE_UPDATE_SNAPSHOT_KEEP, max_file_size=_PRE_UPDATE_SNAPSHOT_MAX_FILE_SIZE, ) - - # Verify live state.db after the snapshot: a concurrent process (antivirus, killed - # gateway, Windows filter driver) can corrupt it and we'd otherwise exit 0 silently. - if snapshot_id: - _src_path = _get_home() / "state.db" - if _src_path.exists(): - _integrity = verify_sqlite_integrity( - _src_path, - check_header=True, - run_pragma=True, - max_bytes=_PRE_UPDATE_SNAPSHOT_MAX_FILE_SIZE, - ) - if not _integrity.get("valid"): - _msg = _integrity.get("message", "unknown error") - print(f" ⚠ state.db integrity check FAILED after snapshot: {_msg}") - _snap_root = _quick_snapshot_root(_get_home()) - _snap_state = _snap_root / snapshot_id / "state.db" - if _snap_state.exists(): - _snap_ok = verify_sqlite_integrity( - _snap_state, check_header=True, run_pragma=True - ) - if _snap_ok.get("valid"): - print(" ✓ Snapshot copy is valid — continuing update.") - print(" If state.db is lost after update it will be auto-restored.") - else: - print( - " ✗ Snapshot copy ALSO failed integrity — " - "the source was already corrupted before the backup." - ) - else: - print(" ⚠ Snapshot does not contain state.db (was skipped or too large).") - print() - if snapshot_id: - print(f"◆ Pre-update snapshot: {snapshot_id}") - - # The code swap + fleet restart touch EVERY profile, so each gets the same snapshot - # under its own state-snapshots/. Best-effort per profile. - with _best_effort('Sibling profile snapshots failed: %s'): - from hermes_cli.backup import create_pre_update_snapshots_all_profiles - - _sibling_snaps = create_pre_update_snapshots_all_profiles( - keep=_PRE_UPDATE_SNAPSHOT_KEEP, - max_file_size=_PRE_UPDATE_SNAPSHOT_MAX_FILE_SIZE, + if _sibling_snaps: + print(f"◆ Sibling profile snapshot(s): " + ", ".join(sorted(_sibling_snaps))) + _record_update_step( + "sibling_profile_snapshots", + True, + ", ".join(f"{k}={v}" for k, v in sorted(_sibling_snaps.items())), ) - if _sibling_snaps: - print(f"◆ Sibling profile snapshot(s): " + ", ".join(sorted(_sibling_snaps))) - _record_update_step( - "sibling_profile_snapshots", - True, - ", ".join( - f"{k}={v}" for k, v in sorted(_sibling_snaps.items()) - ), - ) - import hermes_cli.update_cmd_config as _cfg + import hermes_cli.update_cmd_config as _cfg - # The reader lives in update_cmd_config; write ITS module global, not ours. - _cfg._LAST_SIBLING_SNAPSHOTS = _sibling_snaps + # The reader lives in update_cmd_config; write ITS module global, not ours. + _cfg._LAST_SIBLING_SNAPSHOTS = _sibling_snaps + return snapshot_id - if mode != "full": - if snapshot_id: - print() - return snapshot_id +def _run_full_backup() -> None: + """Zip HERMES_HOME under ``backups/`` (restorable via ``hermes import``). Never raises.""" try: from hermes_cli.backup import create_pre_update_backup except Exception as exc: print(f"⚠ Pre-update backup: could not load backup module ({exc}); continuing update.") print() - return snapshot_id + return try: - from hermes_cli.config import load_config - - _keep = (load_config() or {}).get("updates", {}).get("backup_keep", 5) + _keep = _load_updates_cfg().get("backup_keep", 5) except Exception: _keep = 5 @@ -868,14 +808,13 @@ def _run_pre_update_backup(args) -> Optional[str]: print(f" ⚠ Backup failed: {exc}") print(" Continuing with update.") print() - return snapshot_id - + return elapsed = _time.monotonic() - t0 if out_path is None: print(" ⚠ Backup skipped (no files found or write failed); continuing update.") print() - return snapshot_id + return try: size_bytes = out_path.stat().st_size @@ -884,24 +823,46 @@ def _run_pre_update_backup(args) -> Optional[str]: from hermes_cli.sizefmt import format_bytes - size_str = format_bytes(size_bytes) - # display_hermes_home so the user sees ~/.hermes/... try: from hermes_constants import get_hermes_home, display_hermes_home - home = get_hermes_home() - try: - display_path = f"{display_hermes_home()}/{out_path.relative_to(home)}" - except ValueError: - display_path = str(out_path) + display_path = f"{display_hermes_home()}/{out_path.relative_to(get_hermes_home())}" except Exception: display_path = str(out_path) - print(f" Saved: {display_path} ({size_str}, {elapsed:.1f}s)") + print(f" Saved: {display_path} ({format_bytes(size_bytes)}, {elapsed:.1f}s)") print(f" Restore: hermes import {out_path}") print(" Disable: set updates.pre_update_backup: quick (or off) in config.yaml") print() + + +def _run_pre_update_backup(args) -> Optional[str]: + """Run the pre-update backup; return the quick-snapshot id (None when off/failed). Never raises. + + ``off`` — nothing. ``quick`` (default) — snapshot of critical small files under + ``state-snapshots/``, files over 1 GiB skipped so a bloated state.db can't stall the update. + ``full`` — quick snapshot PLUS a zip of HERMES_HOME under ``backups/`` (``hermes import``). + """ + mode = _resolve_pre_update_backup_mode(args) + + if mode == "off": + if getattr(args, "no_backup", False): + print("◆ Pre-update backup: skipped (--no-backup)") + print() + # Config-level off is silent: the user opted out. + return None + + snapshot_id = None + with _best_effort('Pre-update snapshot failed: %s'): + snapshot_id = _run_quick_snapshots() + + if mode != "full": + if snapshot_id: + print() + return snapshot_id + + _run_full_backup() return snapshot_id @@ -916,15 +877,25 @@ def _sweep_bytecode_after_update(branch: str) -> None: _m()._refresh_bootstrap_cache_scripts(branch) +def _profile_skill_sync_status(r) -> str: + if r and r.get("skipped_opt_out"): + return "opted out (--no-skills)" + if not r: + return "sync failed" + parts = [] + for key, fmt in (("copied", "+{} new"), ("updated", "↑{} updated"), ("user_modified", "~{} user-modified")): + count = len(r.get(key, [])) + if count: + parts.append(fmt.format(count)) + return ", ".join(parts) if parts else "up to date" + + def _sync_profiles_after_update() -> None: """Best-effort per-profile syncs: bundled skills, ``.env`` backfill, Honcho profiles.""" # All profiles incl. the active one: seed_profile_skills() subprocesses with an explicit # HERMES_HOME, so sync_skills()'s module-level HERMES_HOME cache can't skew it. with suppress(Exception): - from hermes_cli.profiles import ( - list_profiles, - seed_profile_skills, - ) + from hermes_cli.profiles import list_profiles, seed_profile_skills all_profiles = list_profiles() if all_profiles: @@ -932,24 +903,7 @@ def _sync_profiles_after_update() -> None: print("→ Syncing bundled skills to all profiles...") for p in all_profiles: try: - r = seed_profile_skills(p.path, quiet=True) - if r and r.get("skipped_opt_out"): - status = "opted out (--no-skills)" - elif r: - copied = len(r.get("copied", [])) - updated = len(r.get("updated", [])) - modified = len(r.get("user_modified", [])) - parts = [] - if copied: - parts.append(f"+{copied} new") - if updated: - parts.append(f"↑{updated} updated") - if modified: - parts.append(f"~{modified} user-modified") - status = ", ".join(parts) if parts else "up to date" - else: - status = "sync failed" - print(f" {p.name}: {status}") + print(f" {p.name}: {_profile_skill_sync_status(seed_profile_skills(p.path, quiet=True))}") except Exception as pe: print(f" {p.name}: error ({pe})") @@ -974,62 +928,56 @@ def _sync_profiles_after_update() -> None: print(f"\n-> Honcho: synced {synced} profile(s)") +def _refresh_cua_driver_after_update() -> None: + """cua-driver refresh, no-op unless on PATH; tied to update for a predictable cadence + without a per-launch GitHub API call.""" + refresh_cua_driver = True + with _best_effort('Could not read updates.refresh_cua_driver: %s'): + refresh_cua_driver = bool(_load_updates_cfg().get("refresh_cua_driver", True)) + + if ( + refresh_cua_driver + and sys.platform in ("darwin", "win32", "linux") + and shutil.which("cua-driver") + ): + from hermes_cli.tools_config import install_cua_driver + + print() + print("→ Refreshing cua-driver (Computer Use)...") + # require_confirmed_update: install only when check-update positively reports a + # newer release (update must stay fast; `computer-use install --upgrade` forces). + # Windows defers even confirmed updates (installer may need console/UAC consent). + install_cua_driver( + upgrade=True, + require_confirmed_update=True, + show_installer_progress=False, + ) + + def _print_post_update_notices_and_self_heals() -> None: """Best-effort notices (FTS optimize, curator) and self-heals (FHS PATH, ACP launcher, Windows bin launchers, cua-driver refresh) that run after the summary.""" from hermes_cli.update_cmd import _m, _print_curator_first_run_notice, _print_curator_recent_run_notice - # v23 FTS layout is opt-in (existing indexes untouched); surface the command here. - with _best_effort('FTS optimize notice failed: %s'): - _print_fts_optimize_available_notice() - - with _best_effort('Curator first-run notice failed: %s'): - _print_curator_first_run_notice() - - with _best_effort('Curator recent-run notice failed: %s'): - _print_curator_recent_run_notice() - - with _best_effort('FHS PATH guard check failed: %s'): - _ensure_fhs_path_guard() - - with _best_effort('hermes-acp launcher self-heal failed: %s'): - _ensure_acp_launcher() - - # Windows launchers into the managed bin dir: in-checkout launchers were swept by the - # autostash (--include-untracked) and updates never run install.ps1. No-op on POSIX. - with _best_effort('Windows bin launcher migration failed: %s'): + def _migrate_windows_bin_path() -> None: + # Windows launchers into the managed bin dir: in-checkout launchers were swept by the + # autostash (--include-untracked) and updates never run install.ps1. No-op on POSIX. from hermes_cli._install_repair import migrate_windows_bin_path migrate_windows_bin_path(_m().PROJECT_ROOT) - # cua-driver refresh, no-op unless on PATH; tied to update for a predictable cadence - # without a per-launch GitHub API call. - with _best_effort('cua-driver refresh failed: %s'): - refresh_cua_driver = True - with _best_effort('Could not read updates.refresh_cua_driver: %s'): - from hermes_cli.config import load_config - - _update_cfg = (load_config() or {}).get("updates", {}) - if isinstance(_update_cfg, dict): - refresh_cua_driver = bool(_update_cfg.get("refresh_cua_driver", True)) - - if ( - refresh_cua_driver - and sys.platform in ("darwin", "win32", "linux") - and shutil.which("cua-driver") - ): - from hermes_cli.tools_config import install_cua_driver - - print() - print("→ Refreshing cua-driver (Computer Use)...") - # require_confirmed_update: install only when check-update positively reports a - # newer release (update must stay fast; `computer-use install --upgrade` forces). - # Windows defers even confirmed updates (installer may need console/UAC consent). - install_cua_driver( - upgrade=True, - require_confirmed_update=True, - show_installer_progress=False, - ) + for message, step in ( + # v23 FTS layout is opt-in (existing indexes untouched); surface the command here. + ('FTS optimize notice failed: %s', _print_fts_optimize_available_notice), + ('Curator first-run notice failed: %s', _print_curator_first_run_notice), + ('Curator recent-run notice failed: %s', _print_curator_recent_run_notice), + ('FHS PATH guard check failed: %s', _ensure_fhs_path_guard), + ('hermes-acp launcher self-heal failed: %s', _ensure_acp_launcher), + ('Windows bin launcher migration failed: %s', _migrate_windows_bin_path), + ('cua-driver refresh failed: %s', _refresh_cua_driver_after_update), + ): + with _best_effort(message): + step() def _run_post_update_maintenance( diff --git a/tests/hermes_cli/test_update_fleet_check_fail_closed.py b/tests/hermes_cli/test_update_fleet_check_fail_closed.py index 53f1a7af20..a0ec5a1cd1 100644 --- a/tests/hermes_cli/test_update_fleet_check_fail_closed.py +++ b/tests/hermes_cli/test_update_fleet_check_fail_closed.py @@ -138,8 +138,14 @@ class TestCallSiteWiring: src = self._impl_source() assert "_fleet_rows_expected = _m()._fleet_probe_expected_runtimes(" in src # The 2.0s settle window must key on the cross-platform signal, so a - # resumed Windows gateway gets its settle window too (#93406). - assert "if _fleet_rows_expected:\n" in src + # resumed Windows gateway gets its settle window too (#93406). The + # settle loop lives in _collect_fleet_snapshot, gated on that signal. + assert "_fleet_snapshot = _collect_fleet_snapshot(restart, _fleet_rows_expected)" in src + from hermes_cli import update_cmd_fleet + + snap_src = inspect.getsource(update_cmd_fleet._collect_fleet_snapshot) + assert "if not rows_expected:\n" in snap_src + assert "_time.sleep(2.0)" in snap_src assert "if restarted_services or killed_pids:\n _time.sleep" not in src def test_zero_row_guard_gated_on_expected_runtimes(self):