"""Gateway fleet restart + post-update verification for ``hermes update``. Split out of ``hermes_cli/update_cmd.py``; every name is re-imported there so ``hermes_cli.update_cmd.`` keeps resolving/monkeypatching. Origin helpers are imported lazily inside each function (no import cycle; test patches stay effective). """ import logging from contextlib import suppress import os import subprocess import sys import time as _time from dataclasses import dataclass, field from pathlib import Path from hermes_cli.update_cmd_common import _best_effort from hermes_cli.update_inventory import _gateway_service_matches_profile # Log-record parity with the origin module. logger = logging.getLogger("hermes_cli.update_cmd") # 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. # The existing ``.update-incomplete`` / ``.lazy-refresh-incomplete`` markers gate dependency/venv repair; # this one is the fleet-restart obligation after a git pull that advanced HEAD (#95294). _FLEET_RESTART_PENDING_NAME = "fleet_restart_pending" _FRESH_RESTART_SUPERVISORS = frozenset({"systemd", "launchd", "service", "s6"}) # A supervisor can report a restarted unit active before the gateway finishes its # bootstrap and publishes ``gateway_state.json``. Keep the readiness poll bounded, # but allow the default systemd startup budget plus status-publication slack. _FLEET_PROBE_SETTLE_TIMEOUT_SECONDS = 120.0 _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 path = get_hermes_home() / ".update_exit_code" with suppress(OSError): path.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.""" from hermes_cli.update_cmd import get_hermes_home return get_hermes_home() / _FLEET_RESTART_PENDING_NAME def _write_fleet_restart_pending_marker(*, expected_sha: str = "") -> None: """Drop the pull→restart obligation breadcrumb. Never raises.""" from hermes_cli.update_cmd import _m path = _fleet_restart_pending_marker_path() if _m()._pytest_owns_live_checkout(path.parent): logger.debug("Skipping fleet-restart-pending marker under pytest (live checkout)") return try: lines = [f"started={_time.time()}", f"pid={os.getpid()}"] if expected_sha: lines.append(f"expected_sha={expected_sha}") path.write_text("\n".join(lines) + "\n", encoding="utf-8") except OSError as exc: logger.debug("Could not write fleet-restart-pending marker: %s", exc) def _clear_fleet_restart_pending_marker() -> None: """Remove the pull→restart obligation breadcrumb. Never raises.""" from hermes_cli.update_cmd import _m _m()._clear_marker_file(_fleet_restart_pending_marker_path(), label="fleet-restart-pending") def _current_checkout_sha() -> str | None: """Current on-disk checkout HEAD, or None if it cannot be resolved.""" from hermes_cli.update_cmd import _capture_head_sha, _m try: from hermes_cli.build_info import get_code_identity sha = (get_code_identity(refresh=True) or {}).get("sha") return str(sha) if sha else None except Exception: return _capture_head_sha(["git"], _m().PROJECT_ROOT) def _receipt_looks_unfinished(receipt: dict) -> bool: """True when *receipt* is from an update that did not finish cleanly. The command boundary stamps a ``stop_reason`` on every receipt, including clean ones (``completed at command boundary``, ``sys.exit(0)``); it must not make a successful receipt look unfinished, or the next ``hermes update`` retriggers ``fleet_restart_pending`` from pre-pull plan SHAs (#98022). """ exit_code = receipt.get("exit_code") outcome = receipt.get("outcome") if exit_code not in (0, None) or 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 # A stop_reason alone (update_contract refusals: outcome="refused", no exit_code) # counts only when nothing else vouched for success. succeeded = exit_code == 0 or outcome == "success" return bool(receipt.get("stop_reason")) and not succeeded def _receipt_reports_stale_runtime(expected_sha: str | None = None) -> bool: """True when ``update_receipts/latest.json`` records a runtime SHA skew. Prefer the post-restart ``fleet`` matrix. ``plan.runtimes[].code_sha`` is captured *before* the pull, so a finished update's plan always looks stale and must not retrigger a restart; consult it only for an unfinished receipt. See #95294. """ from hermes_cli.update_cmd import _current_checkout_sha try: from hermes_cli.update_receipt import read_latest_receipt receipt = read_latest_receipt() except Exception: receipt = None if not isinstance(receipt, dict): return False expected_sha = expected_sha or _current_checkout_sha() if not expected_sha: return False def _sha_mismatch(code_sha) -> bool: return bool(code_sha) and str(code_sha) != str(expected_sha) fleet = receipt.get("fleet") if isinstance(fleet, list) and fleet: 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 return any( isinstance(runtime, dict) and _sha_mismatch(runtime.get("code_sha")) for runtime in plan.get("runtimes") or [] ) def _receipt_owed_gateways() -> set[tuple[str, str]] | None: """``(kind, profile)`` identities ``latest.json`` owes a current successor. Empty when the receipt records no runtimes; ``None`` when any recorded runtime is one the gateway matrix cannot vouch for (serve/dashboard, unknown profile). """ from hermes_cli.update_receipt import read_latest_receipt receipt = read_latest_receipt() or {} plan = receipt.get("plan") or {} entries: list[tuple[object, str | None]] = [(entry, None) for entry in plan.get("runtimes") or []] entries.extend((entry, "gateway") for entry in receipt.get("fleet") or []) owed: set[tuple[str, str]] = set() for entry, default_kind in entries: if not isinstance(entry, dict): return None kind = entry.get("kind", default_kind) profile = entry.get("profile") if kind != "gateway" or not profile or profile == "unknown": return None owed.add((kind, profile)) return owed def _live_fleet_covers_receipt(expected_sha: str | None) -> bool: """Require current successors for every recorded runtime, not just any live row. A PID changes on restart; the stable identity is (runtime kind, profile). The gateway matrix cannot vouch for serve/dashboard or unidentified runtimes. Keep the historical receipt intact: a manual restart is not a successful update. """ if not expected_sha: return False from hermes_cli.update_receipt import collect_fleet_versions try: owed = _receipt_owed_gateways() if not owed: return False fleet = collect_fleet_versions() if not fleet or any( row.get("state") != "current" or row.get("code_sha") != expected_sha for row in fleet ): return False return owed <= {("gateway", row.get("profile")) for row in fleet} except Exception as exc: logger.debug("Could not reconcile pending fleet identities: %s", exc) return False def _read_fleet_marker_expected_sha() -> str: """``expected_sha`` recorded in the pending marker ("" when absent/unreadable).""" with suppress(OSError): for line in _fleet_restart_pending_marker_path().read_text(encoding="utf-8").splitlines(): if line.startswith("expected_sha="): return line.split("=", 1)[1].strip() return "" def _marker_only_restart_obsolete() -> bool: """True when the pending marker is a leftover: the fleet already runs the expected code. A supervisor-level restart (``systemctl --user restart``, launchctl, ops scripts) never goes through this module's clear path, so the marker survives a restart that DID bring every live gateway to the pulled code — and every later CLI call then prints the interrupted-update warning forever (false positive). The marker is not an *unknown* obligation: it records its own ``expected_sha``, so it can be checked against the live fleet directly — no receipt required. Hold it to the same evidence bar ``_live_fleet_covers_receipt`` applies to one: at least one row, and every row a ``current`` gateway under a known profile whose ``code_sha`` equals that ``expected_sha``, with the checkout HEAD not moved past the marker. Keep the marker on stale/down rows, on an all-``unknown`` fleet (pre-code-identity gateways cannot prove currency — same conservatism as the silent-failure class #88848/#74973), on a marker with no ``expected_sha``, when a newer pull moved the checkout, and when the probe fails or answers empty. A gateway the restart phase stopped and never brought back yields NO row at startup (no ``pre_restart_pids`` → no ``down`` classification), so rows alone cannot prove the whole fleet is back. When ``latest.json`` names the gateways the update owed, every one of them must also be covered by a current row; the rows-only rule applies only when the receipt names none. """ expected_sha = _read_fleet_marker_expected_sha() if not expected_sha: return False # pre-expected_sha marker: nothing to verify against checkout_sha = _current_checkout_sha() if checkout_sha and checkout_sha != expected_sha: return False # a newer pull moved HEAD; it owns a fresh obligation try: from hermes_cli.update_receipt import collect_fleet_versions fleet = collect_fleet_versions() owed = _receipt_owed_gateways() except Exception as exc: logger.debug("Fleet probe failed; keeping fleet-restart-pending marker: %s", exc) return False if not fleet: return False # probe answered empty: no proof either way for row in fleet: if not isinstance(row, dict): return False profile = row.get("profile") if not profile or profile == "unknown": return False # unidentified runtime: the matrix cannot vouch for it if row.get("state") != "current" or str(row.get("code_sha")) != expected_sha: return False # stale / down / unknown-identity row still owes the restart if owed is None or not owed <= {("gateway", row.get("profile")) for row in fleet}: return False # a gateway the receipt owes is absent (down) or unidentifiable _clear_fleet_restart_pending_marker() logger.debug( "Fleet-restart-pending marker discharged: %d gateway(s) already serve %s", len(fleet), expected_sha[:10], ) return True def _pending_fleet_restart_needed() -> bool: """Reconcile old restart obligations against current, identity-matched gateways.""" from hermes_cli.update_cmd import _current_checkout_sha # The marker has no runtime inventory and may belong to a newer, killed update # than latest.json. An older receipt cannot discharge that unknown obligation. with suppress(OSError): if _fleet_restart_pending_marker_path().is_file(): if _marker_only_restart_obsolete(): return False return True if not _receipt_reports_stale_runtime(): return False return not _live_fleet_covers_receipt(_current_checkout_sha()) def _warn_pending_fleet_restart(*, startup: bool = False) -> None: """Print the specific interrupted-update fleet-restart warning.""" stream = sys.stderr if startup else sys.stdout print("⚠ A previous `hermes update` pulled new code but did not restart running gateways.", file=stream) print(" Gateways may still be serving pre-update modules (mixed sys.modules).", file=stream) if startup: print(" Run `hermes update` or `hermes gateway restart`.", file=stream) def _warn_pending_fleet_restart_on_startup() -> None: """Cheap CLI-startup hint. Never restarts; never raises.""" with suppress(Exception): 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, listings) -> None: """Best-effort ``systemctl restart`` of every hermes-gateway/serve unit.""" answered = set() for scope, scope_cmd, result in listings: answered.add(scope) if result.returncode != 0: failed.append(f"systemd-{scope} (listing failed)") continue def process_unit(svc_name: str, _scope=scope, _cmd=scope_cmd) -> None: manage_cmd = list(_cmd) + ["--no-ask-password"] if _needs_sudo(_scope): manage_cmd = ["sudo", "-n"] + manage_cmd result = _systemctl_reset_and_restart(manage_cmd, svc_name, scope_cmd=_cmd) if result.returncode != 0 or not _wait_for_service_active(_cmd, svc_name): failed.append(svc_name) _for_each_systemd_gateway_unit( result.stdout, process_unit=process_unit, on_unit_timeout=lambda svc_name, exc: failed.append(svc_name), ) # A timeout or missing executable is not an empty scope. failed.extend(f"systemd-{scope} (listing unavailable)" for scope, _ in _SYSTEMD_SCOPES if scope not in answered) def _run_pending_fleet_restart() -> bool: """Catch-up restart for gateways left on pre-update code. Never raises. True when all discovered targets recovered (or none exist); False if incomplete. See #95294. """ from hermes_cli.update_cmd import _m print("→ Restarting gateways left on pre-update code...") with suppress(Exception): _m()._purge_stale_hermes_modules() # Warn if legacy Hermes gateway unit files are still installed. When both hermes.service (from a # pre-rename install) and the current hermes-gateway.service are enabled, they SIGTERM-fight for the # same bot token (see PR #11909). Flagging here means every `hermes update` surfaces the issue until the # user migrates. try: from hermes_cli.gateway import ( find_gateway_pids, is_macos, is_windows, kill_gateway_processes, supports_systemd_services, _wait_for_gateway_exit, ) except Exception as exc: _warn_gateway_restart_phase_aborted(exc, None) return False try: pids = list(find_gateway_pids(all_profiles=True)) except Exception as exc: logger.debug("Pending fleet restart: gateway probe failed: %s", exc) pids = None failed: list = [] try: # Snapshot before stopping: Restart=no units can disappear from list-units on a clean exit. systemd_listings = list(_systemd_gateway_unit_listings()) if supports_systemd_services() else None # Stop old processes before supervisor recovery, never its freshly verified workers. if pids != []: try: leftover = list(find_gateway_pids(all_profiles=True)) except Exception: leftover = list(pids or []) if leftover: with _best_effort('Pending fleet restart: PID stop failed: %s'): kill_gateway_processes(all_profiles=True) _wait_for_gateway_exit(timeout=5.0, force_after=None) # --- Systemd services (Linux) --- Discover all hermes-gateway* units (default + profiles) plus # hermes-serve* units (the Desktop app's backend, #83438). if systemd_listings is not None: _restart_systemd_gateway_units_best_effort(failed, systemd_listings) # --- Launchd services (macOS) --- Restart EVERY ai.hermes.gateway* LaunchAgent, not only the # invoking profile's — parity with the systemd branch above (#41403). Per-label TimeoutExpired # isolation happens inside. if is_macos(): try: _restart_macos_launchd_gateways([], failed, 45.0, require_supervision=True) except Exception as exc: logger.debug("Pending fleet restart: launchd failed: %s", exc) failed.append("launchd") if is_windows(): try: from hermes_cli import gateway_windows if gateway_windows.is_installed(): gateway_windows.restart() except Exception as exc: logger.debug("Pending fleet restart: Windows failed: %s", exc) failed.append("windows-gateway") if failed: _warn_incomplete_gateway_fleet_restart(failed) return False print(" ✓ Pending fleet restart completed.") return True except Exception as exc: try: surviving = list(find_gateway_pids(all_profiles=True)) except Exception: surviving = pids _warn_gateway_restart_phase_aborted(exc, surviving) return False def _defer_fleet_restart_after_update(*, update_complete: bool, resume_incomplete: bool = False) -> None: """Record a deliberately deferred fleet restart and return/exit on outcome. ``hermes update --no-gateway-restart`` (cron running inside the gateway's own cgroup) updated code and dependencies but must not restart the fleet: the SIGUSR1 drain + systemd restart would kill the updater itself. The ``fleet_restart_pending`` marker written before the pull is KEPT so the next normal update (or ``hermes gateway restart``) catches up. Outcome contract (same success/partial meaning as the normal path): a STALE fleet caused only by this deliberate deferral is expected and does NOT make the update partial — exit 0, receipt "success". The receipt is "partial" (and the process exits 1, marker kept) when the update itself did not complete (``update_complete`` False) or the Windows pause/resume reconciliation reported incomplete. The live interpreter still serves pre-update code here, so callers must also skip stale-module purge/reload work: mutating this process's sys.modules graph mid-flight risks breaking the serving gateway. """ print() print("→ Gateway restart skipped (--no-gateway-restart).") print(" Code and dependencies are updated; gateways still serve pre-update code.") print(" Restart them separately: `hermes gateway restart` or a daily-restart cron.") print(" (fleet restart deferred — marker kept for catch-up)") with suppress(Exception): from hermes_cli.update_receipt import record_skip record_skip("gateway_restart", "--no-gateway-restart: deferred, marker kept") partial = (not update_complete) or resume_incomplete with suppress(Exception): from hermes_cli.update_receipt import finalize_update_receipt finalize_update_receipt("partial" if partial else "success") if partial: sys.exit(1) def _apply_pending_fleet_restart_catchup(*, defer: bool = False) -> None: """On an already-up-to-date ``hermes update``, finish a skipped restart. No-op when nothing is pending; exits 1 on incomplete catch-up so automation does not treat the fleet as healthy. ``defer`` (``--no-gateway-restart``) keeps the marker and warns instead: running the restart from inside the gateway's own cgroup would kill the caller. """ from hermes_cli.update_cmd import _run_pending_fleet_restart if not _pending_fleet_restart_needed(): return if defer: print() _warn_pending_fleet_restart() print(" (fleet restart deferred — --no-gateway-restart; marker kept)") print(" Restart separately: `hermes gateway restart` or next non-cron update.") return print() _warn_pending_fleet_restart() print("→ Running the pending fleet restart...") if _run_pending_fleet_restart(): _clear_fleet_restart_pending_marker() return print(" ⚠ Fleet restart incomplete. Recover with: hermes gateway restart") sys.exit(1) def _systemctl(cmd: list, *, timeout: float): """Run a systemctl (or sudo systemctl) invocation, capturing utf-8 text with a timeout.""" return subprocess.run(cmd, capture_output=True, text=True, encoding="utf-8", errors="replace", timeout=timeout) # poll() takes signed 32-bit milliseconds; keep headroom for rounding in communicate(). _SYSTEMCTL_RESTART_TIMEOUT_MAX = (2**31 - 1) // 1000 - 1 def _systemd_restart_timeout(scope_cmd: list, svc_name: str, *, start_only: bool = False) -> float: """Outwait the unit's stop + start budgets, not just the systemctl client. A client timeout does not cancel the manager's queued restart. Unknown or infinite limits use systemd's usual 90s per phase so automation stays bounded. Custom ExecStop chains or EXTEND_TIMEOUT_USEC can still exceed this budget; genuine timeouts must continue through the existing per-unit failure path. """ from gateway.shutdown_forensics import parse_systemd_duration_to_us budgets = {"TimeoutStartUSec": 90.0} if not start_only: budgets["TimeoutStopUSec"] = 90.0 try: show = _systemctl( scope_cmd + ["show", svc_name, "--property=TimeoutStopUSec,TimeoutStartUSec"], timeout=5, ) except (FileNotFoundError, subprocess.TimeoutExpired): return sum(budgets.values()) + 15.0 if show.returncode == 0: for line in (show.stdout or "").splitlines(): key, _, raw = line.partition("=") if key in budgets: # The shared parser returns None for infinity/unrecognized units. try: raw = raw.strip() duration = int(raw) if raw.isascii() and raw.isdigit() else parse_systemd_duration_to_us(raw) if duration is not None and duration > 0: budgets[key] = duration / 1_000_000 except (ValueError, OverflowError): pass return min(sum(budgets.values()) + 15.0, _SYSTEMCTL_RESTART_TIMEOUT_MAX) def _systemctl_reset_and_restart(manage_cmd: list, svc_name: str, *, scope_cmd: list | None = None): """``reset-failed`` then ``restart``: a unit parked in failed state by systemd's own auto-restart can wedge a plain ``restart`` against RestartSec backoff and stay dead.""" # Property reads need no manage-units privileges: narrow sudoers may permit # restart/reset-failed but deny show. Keep the same user/system manager scope. timeout = _systemd_restart_timeout(scope_cmd if scope_cmd is not None else manage_cmd, svc_name) _systemctl(manage_cmd + ["reset-failed", svc_name], timeout=10) return _systemctl(manage_cmd + ["restart", svc_name], timeout=timeout) 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 ( # list-units is already pattern-filtered, but keep the name gate so a stray non-gateway/serve line # cannot enter the restart path. See #83595. 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, *, process_unit, on_unit_timeout) -> None: """Process each hermes-gateway*/hermes-serve* unit from ``systemctl list-units``. ``TimeoutExpired`` from ``process_unit`` is isolated per unit via ``on_unit_timeout`` so one wedged systemctl call cannot abort the rest of the fleet. See #68523. """ for line in (list_units_stdout or "").strip().splitlines(): parts = line.split() if not parts: continue unit = parts[0] if not unit.endswith(".service") or not _is_hermes_gateway_unit(unit): continue svc_name = unit.removesuffix(".service") try: process_unit(svc_name) except subprocess.TimeoutExpired as exc: on_unit_timeout(svc_name, exc) def _service_unit_supports_graceful_sigusr1_restart(svc_name: str) -> bool: """Whether *svc_name* wires SIGUSR1 to a graceful drain-then-restart. Only ``hermes-gateway*`` runs ``gateway/run.py`` (the handler); SIGUSR1 would just kill ``hermes-serve*`` and burn the drain budget, so those go straight to the blunt restart. Same exact/hyphenated shape as ``_for_each_systemd_gateway_unit`` so a near-prefix unit like ``hermes-gatewayd`` is never signalled. See #83438. """ return svc_name == "hermes-gateway" or svc_name.startswith("hermes-gateway-") def _warn_incomplete_gateway_fleet_restart(failed_units: list) -> None: """Print an explicit incomplete-update warning for unrestarted units.""" from hermes_cli.gateway import is_macos if not failed_units: return ordered = list(dict.fromkeys(failed_units)) # de-dup, discovery order print() print("⚠ Update incomplete — some units were not restarted:") for name in ordered: print(f" - {name}") if is_macos(): # A label lands here when launchd wasn't supervising a live process after # the restart — likely deregistered, which `launchctl kickstart` can't revive. # See #88848. print(" Listed services may be deregistered from launchd, or still") print(" running pre-update code (mixed sys.modules). Recover with:") print(" hermes gateway status") print(" launchctl list | grep