refactor(hermes_cli/update): share profile-home + socket-identity probes between receipt and inventory; fold refusal builder; collapse try/except-pass
This commit is contained in:
@@ -1,4 +1,4 @@
|
||||
"""Image-managed install refusal contract (#91277 Phase 3).
|
||||
"""Image-managed install refusal contract.
|
||||
|
||||
A refusal prints the real update command for the deployment kind, records a ``refused`` receipt (so
|
||||
fleet tooling sees "this install cannot self-update, use <command>" instead of a silent non-update),
|
||||
@@ -10,7 +10,7 @@ from __future__ import annotations
|
||||
import logging
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
from typing import Optional
|
||||
from typing import Callable, Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -24,6 +24,18 @@ class UpdateRefusal:
|
||||
update_command: str # the one-line remediation command
|
||||
|
||||
|
||||
def _refusal(code: str, method: str, message: Optional[Callable[[str], str]] = None) -> UpdateRefusal:
|
||||
"""Refusal for ``method``: ``message(command)`` if given, else docker's full message / the bare command."""
|
||||
from hermes_cli.config import format_docker_update_message, recommended_update_command_for_method
|
||||
|
||||
command = recommended_update_command_for_method(method)
|
||||
if message is not None:
|
||||
text = message(command)
|
||||
else:
|
||||
text = format_docker_update_message() if method == "docker" else command
|
||||
return UpdateRefusal(code=code, message=text, update_command=command)
|
||||
|
||||
|
||||
def evaluate_update_admission(project_root: Path) -> Optional[UpdateRefusal]:
|
||||
"""Return an :class:`UpdateRefusal` when in-place update must not run.
|
||||
|
||||
@@ -36,64 +48,28 @@ def evaluate_update_admission(project_root: Path) -> Optional[UpdateRefusal]:
|
||||
|
||||
provenance = read_image_provenance()
|
||||
if provenance is not None:
|
||||
from hermes_cli.config import (
|
||||
format_docker_update_message,
|
||||
recommended_update_command_for_method,
|
||||
)
|
||||
|
||||
if not provenance.valid:
|
||||
# Present but malformed: still image-managed — an integrity
|
||||
# defect is never permission to mutate the image in place.
|
||||
command = recommended_update_command_for_method("docker")
|
||||
return UpdateRefusal(
|
||||
code="image-marker-invalid",
|
||||
message=(
|
||||
"✗ This install is image-managed, but its provenance "
|
||||
f"marker is invalid ({provenance.error}).\n"
|
||||
" In-place update is disabled. Update by pulling a "
|
||||
f"new image:\n {command}"
|
||||
),
|
||||
update_command=command,
|
||||
)
|
||||
manager = provenance.manager
|
||||
if manager == "docker":
|
||||
return UpdateRefusal(
|
||||
code="image-marker",
|
||||
message=format_docker_update_message(),
|
||||
update_command=recommended_update_command_for_method("docker"),
|
||||
)
|
||||
command = recommended_update_command_for_method(manager)
|
||||
return UpdateRefusal(
|
||||
code="image-marker",
|
||||
message=command,
|
||||
update_command=command,
|
||||
)
|
||||
# Present but malformed: still image-managed — an integrity defect is never
|
||||
# permission to mutate the image in place.
|
||||
return _refusal("image-marker-invalid", "docker", lambda command: (
|
||||
"✗ This install is image-managed, but its provenance "
|
||||
f"marker is invalid ({provenance.error}).\n"
|
||||
" In-place update is disabled. Update by pulling a "
|
||||
f"new image:\n {command}"
|
||||
))
|
||||
return _refusal("image-marker", provenance.manager)
|
||||
except Exception as exc:
|
||||
logger.debug("Image provenance check failed (using heuristics): %s", exc)
|
||||
|
||||
# Layer 2: pre-existing filesystem heuristics, verbatim semantics.
|
||||
try:
|
||||
from hermes_cli.config import (
|
||||
detect_install_method,
|
||||
format_docker_update_message,
|
||||
is_nix_install_method,
|
||||
recommended_update_command_for_method,
|
||||
)
|
||||
from hermes_cli.config import detect_install_method, is_nix_install_method
|
||||
|
||||
method = detect_install_method(project_root)
|
||||
if method == "docker":
|
||||
return UpdateRefusal(
|
||||
code="docker",
|
||||
message=format_docker_update_message(),
|
||||
update_command=recommended_update_command_for_method("docker"),
|
||||
)
|
||||
return _refusal("docker", method)
|
||||
if is_nix_install_method(method) or method == "apt":
|
||||
command = recommended_update_command_for_method(method)
|
||||
return UpdateRefusal(
|
||||
code=method if method == "apt" else "nix",
|
||||
message=command,
|
||||
update_command=command,
|
||||
)
|
||||
return _refusal(method if method == "apt" else "nix", method)
|
||||
except Exception as exc:
|
||||
logger.debug("Install-method admission check failed: %s", exc)
|
||||
return None
|
||||
@@ -106,18 +82,10 @@ def record_refusal_receipt(refusal: UpdateRefusal) -> None:
|
||||
place, use <command>") instead of a silent nothing. Best-effort; never raises.
|
||||
"""
|
||||
try:
|
||||
from hermes_cli.update_receipt import (
|
||||
begin_update_receipt,
|
||||
finalize_update_receipt,
|
||||
record_step,
|
||||
)
|
||||
from hermes_cli.update_receipt import begin_update_receipt, finalize_update_receipt, record_step
|
||||
|
||||
begin_update_receipt()
|
||||
record_step(
|
||||
"admission",
|
||||
False,
|
||||
f"not updatable in place ({refusal.code}); use: {refusal.update_command}",
|
||||
)
|
||||
record_step("admission", False, f"not updatable in place ({refusal.code}); use: {refusal.update_command}")
|
||||
finalize_update_receipt("refused", stop_reason=refusal.code)
|
||||
except Exception as exc:
|
||||
logger.debug("Could not record refusal receipt: %s", exc)
|
||||
|
||||
+120
-266
@@ -1,21 +1,15 @@
|
||||
"""Runtime inventory + update plan for the fleet-update pipeline (#91277 Phase 2).
|
||||
"""Runtime inventory + update plan for the fleet-update pipeline.
|
||||
|
||||
One read-only pass that answers, BEFORE any mutation: what Hermes runtimes are running on this
|
||||
machine, how is each one deployed, which of them will this update touch, and how will each be
|
||||
restarted?
|
||||
|
||||
The module is deliberately side-effect free — every collector is a probe over primitives that
|
||||
already exist (`find_profile_gateway_processes`, `_get_service_pids`, `gateway_state.json` code
|
||||
stamps from #91283, `detect_install_method`) — so `hermes update --plan` can run on a live fleet
|
||||
with zero risk, and the update receipt can embed the inventory without changing update behavior.
|
||||
One read-only pass answering, BEFORE any mutation: which Hermes runtimes run on this machine, how
|
||||
each is deployed, which ones this update touches, and how each restarts. Every collector is a
|
||||
side-effect-free probe, so ``hermes update --plan`` is safe on a live fleet.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from contextlib import contextmanager
|
||||
from contextlib import contextmanager, suppress
|
||||
from dataclasses import dataclass, field, asdict
|
||||
from pathlib import Path
|
||||
from typing import Any, Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -29,7 +23,7 @@ class RuntimeRecord:
|
||||
profile: str # profile name ("default", ...)
|
||||
pid: Optional[int] = None # live PID when known
|
||||
supervisor: str = "manual" # systemd | launchd | desktop | manual
|
||||
code_sha: Optional[str] = None # stamped running-code sha (#91283)
|
||||
code_sha: Optional[str] = None # stamped running-code sha
|
||||
code_version: Optional[str] = None
|
||||
restart_via: str = "" # human-readable restart mechanism
|
||||
detail: dict = field(default_factory=dict)
|
||||
@@ -51,37 +45,34 @@ class UpdatePlan:
|
||||
runtimes: list = field(default_factory=list) # list[RuntimeRecord]
|
||||
|
||||
def to_dict(self) -> dict[str, Any]:
|
||||
return asdict(self) # recursive: RuntimeRecord entries become dicts, dict entries are copied
|
||||
return asdict(self) # recursive: RuntimeRecord entries become dicts
|
||||
|
||||
|
||||
def _detect_supervisor_for_pid(pid: int, service_pids: set, windows_service_pids: set | None = None) -> str:
|
||||
"""Classify how a live gateway PID is supervised."""
|
||||
if windows_service_pids and pid in windows_service_pids:
|
||||
# SCM-supervised Windows gateway (WinSW/NSSM/sc.exe create): the
|
||||
# update pause machinery stops the SERVICE via sc.exe instead of
|
||||
# killing the child, so #91277 Phase 2 reconciliation must plan it
|
||||
# under its own mechanism id, not "manual".
|
||||
# SCM-supervised Windows gateway (WinSW/NSSM/sc.exe create): the update pause machinery
|
||||
# stops the SERVICE via sc.exe instead of killing the child, so reconciliation must plan
|
||||
# it under its own mechanism id, not "manual".
|
||||
return "windows-service"
|
||||
if pid in service_pids:
|
||||
try:
|
||||
from hermes_cli.gateway import is_macos, supports_systemd_services
|
||||
if pid not in service_pids:
|
||||
return "manual"
|
||||
with suppress(Exception):
|
||||
from hermes_cli.gateway import is_macos, supports_systemd_services
|
||||
|
||||
if supports_systemd_services():
|
||||
return "systemd"
|
||||
if is_macos():
|
||||
return "launchd"
|
||||
except Exception:
|
||||
pass
|
||||
return "service"
|
||||
return "manual"
|
||||
if supports_systemd_services():
|
||||
return "systemd"
|
||||
if is_macos():
|
||||
return "launchd"
|
||||
return "service"
|
||||
|
||||
|
||||
# THE restart policy table: restart execution consumes these ids via match_runtime_outcomes / the
|
||||
# update's restart phase, and the receipt records per-runtime outcomes against them. Display
|
||||
# strings are derived by describe_restart_mechanism — never the other way around.
|
||||
_RESTART_MECHANISMS = {
|
||||
"systemd": "systemd",
|
||||
"launchd": "launchd",
|
||||
"desktop": "desktop",
|
||||
"windows-service": "windows-service",
|
||||
"manual-serve": "respawn-argv",
|
||||
"systemd": "systemd", "launchd": "launchd", "desktop": "desktop",
|
||||
"windows-service": "windows-service", "manual-serve": "respawn-argv",
|
||||
}
|
||||
|
||||
_MECHANISM_DESCRIPTIONS = {
|
||||
@@ -92,61 +83,30 @@ _MECHANISM_DESCRIPTIONS = {
|
||||
"respawn-argv": "stop before code swap, relaunch with recorded launch args",
|
||||
}
|
||||
|
||||
_SERVE_KINDS = ("serve", "dashboard")
|
||||
|
||||
|
||||
def _restart_mechanism(supervisor: str, profile: str) -> str:
|
||||
"""Machine-readable restart mechanism id for a runtime.
|
||||
|
||||
THE policy table (#91277 Phase 2): restart execution consumes these ids via
|
||||
:func:`match_runtime_outcomes` / the update's restart phase, and the receipt records per-runtime
|
||||
outcomes against them. Display strings are derived by :func:`describe_restart_mechanism` — never
|
||||
the other way around.
|
||||
"""
|
||||
"""Machine-readable restart mechanism id for a runtime."""
|
||||
return _RESTART_MECHANISMS.get(supervisor, "manual")
|
||||
|
||||
|
||||
def describe_restart_mechanism(mechanism: str, profile: str) -> str:
|
||||
"""Human-readable description of a restart mechanism id."""
|
||||
described = _MECHANISM_DESCRIPTIONS.get(mechanism)
|
||||
if described is not None:
|
||||
return described
|
||||
if profile != "default":
|
||||
return f"hermes -p {profile} gateway restart"
|
||||
return "hermes gateway restart"
|
||||
|
||||
|
||||
def _runtime(kind: str, profile: str, pid: Optional[int], supervisor: str, **extra: Any) -> RuntimeRecord:
|
||||
"""A :class:`RuntimeRecord` with ``restart_via`` derived from its supervisor."""
|
||||
return RuntimeRecord(
|
||||
kind=kind,
|
||||
profile=profile,
|
||||
pid=pid,
|
||||
supervisor=supervisor,
|
||||
restart_via=_restart_mechanism(supervisor, profile),
|
||||
**extra,
|
||||
return _MECHANISM_DESCRIPTIONS.get(mechanism) or (
|
||||
f"hermes -p {profile} gateway restart" if profile != "default" else "hermes gateway restart"
|
||||
)
|
||||
|
||||
|
||||
def _int_or_none(value: Any) -> Optional[int]:
|
||||
try:
|
||||
return int(value)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _gateway_record(
|
||||
profile: str,
|
||||
pid: int,
|
||||
supervisor: str,
|
||||
code_sha: Any = None,
|
||||
code_version: Any = None,
|
||||
def _runtime(
|
||||
kind: str, profile: str, pid: Optional[int], supervisor: str,
|
||||
code_sha: Any = None, code_version: Any = None, **extra: Any,
|
||||
) -> RuntimeRecord:
|
||||
return _runtime(
|
||||
"gateway",
|
||||
profile,
|
||||
pid,
|
||||
supervisor,
|
||||
code_sha=str(code_sha) if code_sha else None,
|
||||
code_version=code_version,
|
||||
"""A :class:`RuntimeRecord` with ``restart_via`` derived from its supervisor."""
|
||||
return RuntimeRecord(
|
||||
kind=kind, profile=profile, pid=pid, supervisor=supervisor,
|
||||
code_sha=str(code_sha) if code_sha else None, code_version=code_version,
|
||||
restart_via=_restart_mechanism(supervisor, profile), **extra,
|
||||
)
|
||||
|
||||
|
||||
@@ -166,14 +126,8 @@ def collect_runtime_inventory() -> UpdatePlan:
|
||||
The result is embeddable in the update receipt and printable via :func:`print_update_plan`.
|
||||
"""
|
||||
plan = UpdatePlan()
|
||||
|
||||
# --- install shape / deployment kind ---------------------------------
|
||||
with _probe("Install-method probe"):
|
||||
from hermes_cli.config import (
|
||||
detect_install_method,
|
||||
get_managed_system,
|
||||
recommended_update_command_for_method,
|
||||
)
|
||||
from hermes_cli.config import detect_install_method, get_managed_system, recommended_update_command_for_method
|
||||
|
||||
method = detect_install_method()
|
||||
plan.install_method = method
|
||||
@@ -181,11 +135,10 @@ def collect_runtime_inventory() -> UpdatePlan:
|
||||
if managed:
|
||||
plan.install_method = managed
|
||||
plan.updatable_in_place = method in ("git", "unknown") and not managed
|
||||
# Baked image provenance (#91277 Phase 3): when the image marker is
|
||||
# present it is authoritative — a bind-mounted checkout inside a
|
||||
# container can look like `git` to the heuristics while the running
|
||||
# filesystem is actually an immutable image. Fail-closed: an invalid
|
||||
# marker still flips the plan to not-updatable.
|
||||
# Baked image provenance: when the image marker is present it is authoritative — a
|
||||
# bind-mounted checkout inside a container can look like `git` to the heuristics while
|
||||
# the running filesystem is actually an immutable image. Fail-closed: an invalid marker
|
||||
# still flips the plan to not-updatable.
|
||||
with _probe("Image provenance probe"):
|
||||
from hermes_cli.image_provenance import read_image_provenance
|
||||
|
||||
@@ -195,102 +148,64 @@ def collect_runtime_inventory() -> UpdatePlan:
|
||||
if provenance.valid and provenance.manager:
|
||||
plan.install_method = provenance.manager
|
||||
plan.update_mechanism = recommended_update_command_for_method(method)
|
||||
|
||||
# --- expected code identity (pre-pull) --------------------------------
|
||||
with _probe("Code-identity probe"):
|
||||
from hermes_cli.build_info import get_code_identity
|
||||
|
||||
identity = get_code_identity(refresh=True)
|
||||
plan.expected_sha = identity.get("sha")
|
||||
plan.expected_version = identity.get("version")
|
||||
|
||||
# --- profiles ----------------------------------------------------------
|
||||
profile_homes: list[tuple[str, Path]] = []
|
||||
profile_homes = []
|
||||
with _probe("Profile enumeration"):
|
||||
from hermes_cli.profiles import (
|
||||
_get_default_hermes_home,
|
||||
_get_profiles_root,
|
||||
_PROFILE_ID_RE,
|
||||
)
|
||||
from hermes_cli.update_receipt import _profile_homes
|
||||
|
||||
default_home = _get_default_hermes_home()
|
||||
if default_home.is_dir():
|
||||
profile_homes.append(("default", default_home))
|
||||
root = _get_profiles_root()
|
||||
if root.is_dir():
|
||||
for entry in sorted(root.iterdir()):
|
||||
if entry.is_dir() and entry.name != "default" and _PROFILE_ID_RE.match(entry.name):
|
||||
profile_homes.append((entry.name, entry))
|
||||
profile_homes = _profile_homes()
|
||||
plan.profiles = [name for name, _ in profile_homes]
|
||||
|
||||
# --- service-managed PIDs (fleet-wide) ---------------------------------
|
||||
service_pids: set = set()
|
||||
with _probe("Service-PID probe"):
|
||||
from hermes_cli.gateway import _get_service_pids
|
||||
|
||||
service_pids = _get_service_pids(all_profiles=True) or set()
|
||||
|
||||
# --- SCM-supervised gateway PIDs (Windows) ------------------------------
|
||||
# find_windows_gateway_services() maps validated gateway PIDs through
|
||||
# process ancestry to running SCM service PIDs (no-op off Windows). The
|
||||
# update's pause phase stops these via `sc.exe stop` / restarts via
|
||||
# `sc.exe start`, so the plan must carry the matching mechanism id for
|
||||
# the #91277 Phase 2 reconciliation and the fleet check.
|
||||
# Windows SCM services (no-op off Windows): the update's pause phase stops these via `sc.exe
|
||||
# stop` / restarts via `sc.exe start`, so the plan must carry the matching mechanism id.
|
||||
windows_service_pids: set = set()
|
||||
with _probe("Windows SCM service-ownership probe"):
|
||||
from hermes_cli.gateway import find_windows_gateway_services
|
||||
|
||||
windows_service_pids = {int(service.gateway_pid) for service in find_windows_gateway_services()}
|
||||
|
||||
# --- per-profile gateways (PID files + runtime status stamps) ----------
|
||||
def _supervisor(pid: int) -> str:
|
||||
return _detect_supervisor_for_pid(pid, service_pids, windows_service_pids)
|
||||
# Per-profile gateways: control-socket identity first (declared by the process itself,
|
||||
# including supervisor provenance — no argv/PID inference), gateway_state.json fallback.
|
||||
seen_pids: set[int] = set()
|
||||
with _probe("Gateway-state inventory"):
|
||||
from gateway.status import _pid_exists, read_runtime_status
|
||||
from hermes_cli.update_receipt import _socket_identity
|
||||
|
||||
for profile, home in profile_homes:
|
||||
# Prefer the gateway-owned control socket (#92091): identity
|
||||
# declared by the process itself, including its own supervisor
|
||||
# provenance — no argv/PID inference. Scan fallback below.
|
||||
identity = None
|
||||
try:
|
||||
from gateway.control_socket import identify_gateway
|
||||
|
||||
identity = identify_gateway(home)
|
||||
except Exception:
|
||||
identity = None
|
||||
sock_pid = _int_or_none(identity.get("pid")) if identity else None
|
||||
if sock_pid is not None:
|
||||
sock = _socket_identity(home)
|
||||
if sock is not None:
|
||||
sock_pid, identity = sock
|
||||
if sock_pid in seen_pids:
|
||||
# One multiplex gateway can answer identify for
|
||||
# several profile homes — one runtime record per
|
||||
# process, not per home.
|
||||
# One multiplex gateway answers identify for several profile homes — one record per process.
|
||||
continue
|
||||
seen_pids.add(sock_pid)
|
||||
declared = identity.get("supervisor")
|
||||
supervisor = (
|
||||
str(declared)
|
||||
if declared
|
||||
else _detect_supervisor_for_pid(sock_pid, service_pids, windows_service_pids)
|
||||
)
|
||||
plan.runtimes.append(
|
||||
_gateway_record(
|
||||
profile, sock_pid, supervisor, identity.get("code_sha"), identity.get("code_version")
|
||||
)
|
||||
)
|
||||
continue
|
||||
record = read_runtime_status(home / "gateway_state.json") or {}
|
||||
pid = _int_or_none(record.get("pid"))
|
||||
if pid is None or not _pid_exists(pid):
|
||||
continue
|
||||
seen_pids.add(pid)
|
||||
supervisor = str(declared) if declared else _supervisor(sock_pid)
|
||||
record: dict = identity
|
||||
else:
|
||||
record = read_runtime_status(home / "gateway_state.json") or {}
|
||||
try:
|
||||
sock_pid = int(record.get("pid"))
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
if not _pid_exists(sock_pid):
|
||||
continue
|
||||
seen_pids.add(sock_pid)
|
||||
supervisor = _supervisor(sock_pid)
|
||||
plan.runtimes.append(
|
||||
_gateway_record(
|
||||
profile,
|
||||
pid,
|
||||
_detect_supervisor_for_pid(pid, service_pids, windows_service_pids),
|
||||
record.get("code_sha"),
|
||||
record.get("code_version"),
|
||||
)
|
||||
_runtime("gateway", profile, sock_pid, supervisor, record.get("code_sha"), record.get("code_version"))
|
||||
)
|
||||
|
||||
# PID-file mapped gateways not covered by a runtime-status record
|
||||
@@ -301,69 +216,47 @@ def collect_runtime_inventory() -> UpdatePlan:
|
||||
if proc.pid in seen_pids:
|
||||
continue
|
||||
seen_pids.add(proc.pid)
|
||||
plan.runtimes.append(
|
||||
_gateway_record(
|
||||
proc.profile, proc.pid, _detect_supervisor_for_pid(proc.pid, service_pids, windows_service_pids)
|
||||
)
|
||||
)
|
||||
plan.runtimes.append(_runtime("gateway", proc.profile, proc.pid, _supervisor(proc.pid)))
|
||||
|
||||
# Serve/dashboard backends from the spawn ledger (#63206). These are the
|
||||
# runtimes the gateway collectors above can never see: a manually
|
||||
# launched `hermes serve --host <ip>` for a remote Desktop, or a
|
||||
# long-lived `hermes dashboard`. Every serve/dashboard registers itself
|
||||
# (with structured host/port/profile since #63206) at startup, and
|
||||
# ledger_entries() live-verifies (pid, create_time) so PID reuse never
|
||||
# fabricates a row. Desktop-supervised backends are classified by their
|
||||
# recorded spawner still being alive — those restart via the Desktop's
|
||||
# own respawn, not ours.
|
||||
# Serve/dashboard backends from the spawn ledger — runtimes the gateway collectors can never
|
||||
# see (a manual `hermes serve --host <ip>` for a remote Desktop, a long-lived `hermes dashboard`).
|
||||
# ledger_entries() live-verifies (pid, create_time) so PID reuse never fabricates a row. Desktop-
|
||||
# supervised backends (spawner still alive) restart via the Desktop's own respawn, not ours.
|
||||
with _probe("Serve/dashboard ledger inventory"):
|
||||
from hermes_cli.process_identity import ledger_entries, spawner_is_dead
|
||||
|
||||
for entry in ledger_entries():
|
||||
purpose = entry.get("purpose")
|
||||
if purpose not in ("serve", "dashboard"):
|
||||
if purpose not in _SERVE_KINDS:
|
||||
continue
|
||||
pid = entry.get("pid")
|
||||
if not isinstance(pid, int) or pid in seen_pids:
|
||||
continue
|
||||
seen_pids.add(pid)
|
||||
plan.runtimes.append(
|
||||
_runtime(
|
||||
str(purpose),
|
||||
str(entry.get("profile") or "default"),
|
||||
pid,
|
||||
"desktop" if spawner_is_dead(entry) is False else "manual-serve",
|
||||
detail={
|
||||
"argv": entry.get("argv") or "",
|
||||
"host": entry.get("host") or "",
|
||||
"port": entry.get("port"),
|
||||
# Process incarnation, not just the numeric PID: a
|
||||
# post-update survivor probe that compares PIDs alone
|
||||
# calls a NEW serve that reused the number a survivor
|
||||
# (#92145 review).
|
||||
"create_time": entry.get("create_time"),
|
||||
},
|
||||
)
|
||||
)
|
||||
|
||||
# detail.create_time: process incarnation, not just the numeric PID — a post-update
|
||||
# survivor probe comparing PIDs alone calls a NEW serve that reused the number a survivor.
|
||||
plan.runtimes.append(_runtime(
|
||||
str(purpose), str(entry.get("profile") or "default"), pid,
|
||||
"desktop" if spawner_is_dead(entry) is False else "manual-serve",
|
||||
detail={
|
||||
"argv": entry.get("argv") or "", "host": entry.get("host") or "",
|
||||
"port": entry.get("port"), "create_time": entry.get("create_time"),
|
||||
},
|
||||
))
|
||||
return plan
|
||||
|
||||
|
||||
def print_update_plan(plan: UpdatePlan) -> None:
|
||||
"""Human-readable plan — what the update will touch and how."""
|
||||
print("Update plan:")
|
||||
print(f" Install: {plan.install_method}", end="")
|
||||
install = f" Install: {plan.install_method}"
|
||||
if plan.expected_version:
|
||||
print(f" (v{plan.expected_version}", end="")
|
||||
if plan.expected_sha:
|
||||
print(f" @ {plan.expected_sha[:8]}", end="")
|
||||
print(")", end="")
|
||||
print()
|
||||
install += f" (v{plan.expected_version}" + (f" @ {plan.expected_sha[:8]}" if plan.expected_sha else "") + ")"
|
||||
print(install)
|
||||
if not plan.updatable_in_place:
|
||||
print(" ⚠ This install is NOT updatable in place.")
|
||||
print(f" Update via: {plan.update_mechanism}")
|
||||
profiles = ", ".join(plan.profiles) if plan.profiles else "(none found)"
|
||||
print(f" Profiles: {profiles}")
|
||||
print(f" Profiles: {', '.join(plan.profiles) if plan.profiles else '(none found)'}")
|
||||
if not plan.runtimes:
|
||||
print(" Running Hermes services: none detected — code swap only.")
|
||||
return
|
||||
@@ -374,54 +267,17 @@ def print_update_plan(plan: UpdatePlan) -> None:
|
||||
print(f" restart: {describe_restart_mechanism(runtime.restart_via, runtime.profile)}")
|
||||
|
||||
|
||||
_SERVE_KINDS = ("serve", "dashboard")
|
||||
|
||||
|
||||
def _serve_unit_matches_profile(profile: str, unit: object) -> bool:
|
||||
"""Does *unit* name a ``hermes-serve*``/``hermes-dashboard*`` unit for *profile*?
|
||||
|
||||
Serve/dashboard runtimes have their OWN unit vocabulary; the gateway's ``hermes-gateway*`` names
|
||||
never cover them (#100479).
|
||||
"""
|
||||
name = str(unit).removesuffix(".service")
|
||||
if "/" in name:
|
||||
name = name.rsplit("/", 1)[-1]
|
||||
if profile == "default":
|
||||
return name in {"hermes-serve", "hermes-dashboard"}
|
||||
return name in {f"hermes-serve-{profile}", f"hermes-dashboard-{profile}"}
|
||||
|
||||
|
||||
def _serve_runtime_outcome(
|
||||
r: RuntimeRecord,
|
||||
*,
|
||||
killed: set,
|
||||
failed_set: set,
|
||||
restarted_set: set,
|
||||
stale_serves: "set | None",
|
||||
) -> str:
|
||||
"""Outcome for one serve/dashboard runtime — never the gateway's."""
|
||||
if r.pid is not None and r.pid in killed:
|
||||
return "stopped"
|
||||
if any(_serve_unit_matches_profile(r.profile, u) for u in failed_set):
|
||||
return "failed"
|
||||
if stale_serves is not None:
|
||||
# Incarnation-verified: the pre-update process is gone (replaced by
|
||||
# its unit / the dashboard cleanup respawn / the Desktop app) or it
|
||||
# is still alive on pre-update code.
|
||||
return "unaccounted" if r.pid in stale_serves else "restarted"
|
||||
if any(_serve_unit_matches_profile(r.profile, s) for s in restarted_set):
|
||||
return "restarted"
|
||||
return "unaccounted"
|
||||
"""Does *unit* name a ``hermes-serve*``/``hermes-dashboard*`` unit for *profile*? (OWN vocabulary;
|
||||
the gateway's ``hermes-gateway*`` names never cover serve/dashboard runtimes.)"""
|
||||
name = str(unit).removesuffix(".service").rsplit("/", 1)[-1]
|
||||
suffix = "" if profile == "default" else f"-{profile}"
|
||||
return name in {f"hermes-serve{suffix}", f"hermes-dashboard{suffix}"}
|
||||
|
||||
|
||||
def match_runtime_outcomes(
|
||||
plan: "UpdatePlan",
|
||||
*,
|
||||
restarted_services: list,
|
||||
relaunched_profiles: list,
|
||||
externally_supervised_profiles: list,
|
||||
killed_pids: set,
|
||||
failed_units: list,
|
||||
plan: "UpdatePlan", *, restarted_services: list, relaunched_profiles: list,
|
||||
externally_supervised_profiles: list, killed_pids: set, failed_units: list,
|
||||
stale_serve_pids: "set | None" = None,
|
||||
) -> list[dict[str, Any]]:
|
||||
"""Reconcile the plan's runtimes against what the restart phase DID.
|
||||
@@ -444,31 +300,33 @@ def match_runtime_outcomes(
|
||||
stale_serves = {int(p) for p in stale_serve_pids} if stale_serve_pids is not None else None
|
||||
|
||||
def _gateway_names(r: RuntimeRecord, names: set) -> bool:
|
||||
# The bare "hermes-gateway" unit name is gateway-specific: a
|
||||
# serve/dashboard runtime that merely shares the default
|
||||
# profile is a different process the gateway restart never
|
||||
# touched, and must not borrow its outcome (#100479).
|
||||
# The bare "hermes-gateway" unit name is gateway-specific: a serve/dashboard runtime
|
||||
# that merely shares the default profile is a different process the gateway restart
|
||||
# never touched, and must not borrow its outcome.
|
||||
return any(
|
||||
r.profile in name
|
||||
or (
|
||||
r.kind == "gateway"
|
||||
and r.profile == "default"
|
||||
and "hermes-gateway" in name
|
||||
)
|
||||
r.profile in name or (r.kind == "gateway" and r.profile == "default" and "hermes-gateway" in name)
|
||||
for name in names
|
||||
)
|
||||
|
||||
def _serve_outcome(r: RuntimeRecord) -> str:
|
||||
# Serve/dashboard outcome in its OWN vocabulary — never the gateway's.
|
||||
if r.pid is not None and r.pid in killed:
|
||||
return "stopped"
|
||||
if any(_serve_unit_matches_profile(r.profile, u) for u in failed_set):
|
||||
return "failed"
|
||||
if stale_serves is not None:
|
||||
# Incarnation-verified: the pre-update process is gone (replaced by its unit / the
|
||||
# dashboard cleanup respawn / the Desktop app) or it is still alive on pre-update code.
|
||||
return "unaccounted" if r.pid in stale_serves else "restarted"
|
||||
if any(_serve_unit_matches_profile(r.profile, s) for s in restarted_set):
|
||||
return "restarted"
|
||||
return "unaccounted"
|
||||
|
||||
for r in plan.runtimes:
|
||||
if not isinstance(r, RuntimeRecord):
|
||||
continue
|
||||
if r.kind in _SERVE_KINDS:
|
||||
outcome = _serve_runtime_outcome(
|
||||
r,
|
||||
killed=killed,
|
||||
failed_set=failed_set,
|
||||
restarted_set=restarted_set,
|
||||
stale_serves=stale_serves,
|
||||
)
|
||||
outcome = _serve_outcome(r)
|
||||
elif r.profile in relaunched or r.profile in external:
|
||||
outcome = "restarted"
|
||||
elif r.pid is not None and r.pid in killed:
|
||||
@@ -479,13 +337,9 @@ def match_runtime_outcomes(
|
||||
outcome = "restarted"
|
||||
else:
|
||||
outcome = "unaccounted"
|
||||
outcomes.append({
|
||||
"kind": r.kind,
|
||||
"profile": r.profile,
|
||||
"pid": r.pid,
|
||||
"mechanism": r.restart_via,
|
||||
"outcome": outcome,
|
||||
})
|
||||
outcomes.append(
|
||||
{"kind": r.kind, "profile": r.profile, "pid": r.pid, "mechanism": r.restart_via, "outcome": outcome}
|
||||
)
|
||||
except Exception as exc:
|
||||
logger.debug("Runtime-outcome reconciliation failed: %s", exc)
|
||||
return outcomes
|
||||
@@ -510,8 +364,8 @@ def report_unaccounted_runtimes(outcomes: list[dict[str, Any]]) -> bool:
|
||||
print(" hermes gateway restart # active profile")
|
||||
print(" hermes -p <profile> gateway restart # named profile")
|
||||
if any(o.get("kind") in _SERVE_KINDS for o in missed):
|
||||
# A serve/dashboard is not reachable by any `gateway restart`
|
||||
# command (#100479): name the process, not the wrong verb.
|
||||
# A serve/dashboard is not reachable by any `gateway restart` command: name the process,
|
||||
# not the wrong verb.
|
||||
print(" systemctl --user restart hermes-serve.service # unit-managed serve")
|
||||
print(" relaunch `hermes serve` / `hermes dashboard` / the Desktop app")
|
||||
return True
|
||||
|
||||
+124
-229
@@ -1,32 +1,27 @@
|
||||
"""Structured update receipts + post-update fleet version verification.
|
||||
|
||||
Phase 1 of the fleet-update reliability plan (#91277): the updater must *prove* its outcome instead
|
||||
of assuming it.
|
||||
|
||||
Two additive capabilities, both designed so a failure inside them can never break an update (every
|
||||
public entry point is exception-swallowing):
|
||||
The updater must *prove* its outcome instead of assuming it. Every public entry point is
|
||||
exception-swallowing so a failure inside receipts can never break an update.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
from contextlib import suppress
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any, Optional
|
||||
|
||||
import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_RECEIPT_DIR_NAME = "update_receipts"
|
||||
_RECEIPT_KEEP = 20 # keep the last N receipts per profile home
|
||||
|
||||
# Module-level current receipt. ``hermes update`` is a single-threaded CLI
|
||||
# command; a module singleton lets the 7k-line updater record steps from
|
||||
# any depth without threading a handle through every helper.
|
||||
# ``hermes update`` is a single-threaded CLI command; a module singleton lets the 7k-line updater
|
||||
# record steps from any depth without threading a handle through every helper.
|
||||
_current: Optional["UpdateReceipt"] = None
|
||||
|
||||
|
||||
@@ -51,89 +46,60 @@ class UpdateReceipt:
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.data: dict[str, Any] = {
|
||||
"schema": 1,
|
||||
"started_at": _utc_now_iso(),
|
||||
"finished_at": None,
|
||||
"argv": list(sys.argv),
|
||||
"pid": os.getpid(),
|
||||
"schema": 1, "started_at": _utc_now_iso(), "finished_at": None,
|
||||
"argv": list(sys.argv), "pid": os.getpid(),
|
||||
"outcome": "running", # running | success | partial | failed
|
||||
"pre_update": {},
|
||||
"post_update": {},
|
||||
"steps": [],
|
||||
"skips": [],
|
||||
"gateway_restart": {},
|
||||
"fleet": [],
|
||||
"pre_update": {}, "post_update": {},
|
||||
"steps": [], "skips": [], "gateway_restart": {}, "fleet": [],
|
||||
}
|
||||
try:
|
||||
with suppress(Exception):
|
||||
from hermes_cli.build_info import get_code_identity
|
||||
|
||||
self.data["pre_update"] = get_code_identity()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# -- recording ---------------------------------------------------------
|
||||
def step(self, name: str, ok: bool, detail: str = "") -> None:
|
||||
self.data["steps"].append(
|
||||
{"name": name, "ok": bool(ok), "detail": detail, "at": _utc_now_iso()}
|
||||
)
|
||||
self.data["steps"].append({"name": name, "ok": bool(ok), "detail": detail, "at": _utc_now_iso()})
|
||||
|
||||
def skip(self, name: str, reason: str) -> None:
|
||||
self.data["skips"].append(
|
||||
{"name": name, "reason": reason, "at": _utc_now_iso()}
|
||||
)
|
||||
self.data["skips"].append({"name": name, "reason": reason, "at": _utc_now_iso()})
|
||||
|
||||
def gateway_restart_result(
|
||||
self,
|
||||
*,
|
||||
restarted_services: list | None = None,
|
||||
relaunched_profiles: list | None = None,
|
||||
externally_supervised_profiles: list | None = None,
|
||||
killed_pids: list | None = None,
|
||||
failed_units: list | None = None,
|
||||
incomplete: bool = False,
|
||||
phase_error: str = "",
|
||||
self, *, restarted_services: list | None = None, relaunched_profiles: list | None = None,
|
||||
externally_supervised_profiles: list | None = None, killed_pids: list | None = None,
|
||||
failed_units: list | None = None, incomplete: bool = False, phase_error: str = "",
|
||||
fresh_recovery: dict[str, Any] | None = None,
|
||||
) -> None:
|
||||
result: dict[str, Any] = {
|
||||
"restarted_services": list(restarted_services or []),
|
||||
"relaunched_profiles": list(relaunched_profiles or []),
|
||||
"externally_supervised_profiles": list(
|
||||
externally_supervised_profiles or []
|
||||
),
|
||||
"externally_supervised_profiles": list(externally_supervised_profiles or []),
|
||||
"killed_pids": [int(p) for p in (killed_pids or [])],
|
||||
"failed_units": [str(u) for u in (failed_units or [])],
|
||||
"incomplete": bool(incomplete),
|
||||
"phase_error": phase_error,
|
||||
}
|
||||
if fresh_recovery is not None:
|
||||
# Conservative outcome vocabulary: "verified" is the only bucket
|
||||
# allowed to claim supervisor coverage; "relaunch_attempted" means
|
||||
# the relaunch exited 0 without independent supervisor
|
||||
# observation. "skipped" preserves runtimes (manual gateways,
|
||||
# Conservative outcome vocabulary: "verified" is the only bucket allowed to claim
|
||||
# supervisor coverage; "relaunch_attempted" means the relaunch exited 0 without
|
||||
# independent supervisor observation. "skipped" preserves runtimes (manual gateways,
|
||||
# serve/dashboard entries) the pass deliberately did not touch.
|
||||
persisted: dict[str, Any] = {
|
||||
key: [str(profile) for profile in fresh_recovery.get(key, [])]
|
||||
for key in ("requested", "verified", "relaunch_attempted", "failed")
|
||||
}
|
||||
persisted["skipped"] = _str_records(
|
||||
fresh_recovery.get("skipped", []),
|
||||
("profile", "kind", "supervisor", "reason"),
|
||||
fresh_recovery.get("skipped", []), ("profile", "kind", "supervisor", "reason")
|
||||
)
|
||||
# Serve/dashboard coverage (#92145). ``hermes serve`` hosts
|
||||
# tui_gateway and is not a gateway profile, so neither the
|
||||
# per-profile buckets above nor the fleet-version matrix can
|
||||
# describe it. Persist its unit outcomes and any process that
|
||||
# survived on the pre-update generation, or the receipt keeps
|
||||
# claiming a clean recovery the operator's box contradicts.
|
||||
# ``hermes serve`` hosts tui_gateway and is not a gateway profile, so neither the
|
||||
# per-profile buckets above nor the fleet-version matrix can describe it. Persist its
|
||||
# unit outcomes and any process that survived on the pre-update generation, or the
|
||||
# receipt keeps claiming a clean recovery the operator's box contradicts.
|
||||
serve_units = fresh_recovery.get("serve_units") or {}
|
||||
persisted["serve_units"] = {
|
||||
key: [str(unit) for unit in (serve_units.get(key) or [])]
|
||||
for key in ("verified", "failed")
|
||||
key: [str(unit) for unit in (serve_units.get(key) or [])] for key in ("verified", "failed")
|
||||
}
|
||||
persisted["stale_runtimes"] = _str_records(
|
||||
fresh_recovery.get("stale_runtimes", []),
|
||||
("kind", "profile", "supervisor"),
|
||||
pid=True,
|
||||
fresh_recovery.get("stale_runtimes", []), ("kind", "profile", "supervisor"), pid=True
|
||||
)
|
||||
result["fresh_recovery"] = persisted
|
||||
self.data["gateway_restart"] = result
|
||||
@@ -141,18 +107,16 @@ class UpdateReceipt:
|
||||
def finalize(self, outcome: str) -> None:
|
||||
self.data["outcome"] = outcome
|
||||
self.data["finished_at"] = _utc_now_iso()
|
||||
try:
|
||||
with suppress(Exception):
|
||||
from hermes_cli.build_info import get_code_identity
|
||||
|
||||
self.data["post_update"] = get_code_identity(refresh=True)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def _receipt_dir() -> Path:
|
||||
from hermes_cli.config import get_hermes_home
|
||||
|
||||
return get_hermes_home() / "logs" / _RECEIPT_DIR_NAME
|
||||
return get_hermes_home() / "logs" / "update_receipts"
|
||||
|
||||
|
||||
def begin_update_receipt() -> None:
|
||||
@@ -189,14 +153,11 @@ def record_gateway_restart(**kwargs: Any) -> None:
|
||||
_record("gateway_restart_result", "gateway restart result", **kwargs)
|
||||
|
||||
|
||||
def finalize_update_receipt(
|
||||
outcome: str, fleet: list | None = None, stop_reason: str = ""
|
||||
) -> Optional[Path]:
|
||||
"""Finalize + persist the receipt. Returns the written path or None.
|
||||
def finalize_update_receipt(outcome: str, fleet: list | None = None, stop_reason: str = "") -> Optional[Path]:
|
||||
"""Finalize + persist the receipt (``success``/``partial``/``failed``/``refused``); path or None.
|
||||
|
||||
``outcome`` is one of ``success`` / ``partial`` / ``failed`` / ``refused``. Exactly-once by
|
||||
construction: the module singleton is popped first, so a second call (e.g. the command-boundary
|
||||
safety net after an inner path already finalized) is a no-op returning None.
|
||||
Exactly-once by construction: the module singleton is popped first, so a second call (e.g. the
|
||||
command-boundary safety net after an inner path already finalized) is a no-op returning None.
|
||||
"""
|
||||
global _current
|
||||
receipt = _current
|
||||
@@ -211,15 +172,12 @@ def finalize_update_receipt(
|
||||
receipt.data["fleet"] = fleet
|
||||
directory = _receipt_dir()
|
||||
directory.mkdir(parents=True, exist_ok=True)
|
||||
stamp = time.strftime("%Y%m%d_%H%M%S")
|
||||
path = directory / f"update_{stamp}_{os.getpid()}.json"
|
||||
path = directory / f"update_{time.strftime('%Y%m%d_%H%M%S')}_{os.getpid()}.json"
|
||||
body = json.dumps(receipt.data, indent=2, default=str)
|
||||
path.write_text(body, encoding="utf-8")
|
||||
# Stable pointer for the dashboard/desktop: latest receipt.
|
||||
try:
|
||||
with suppress(OSError):
|
||||
(directory / "latest.json").write_text(body, encoding="utf-8")
|
||||
except OSError:
|
||||
pass
|
||||
_prune_old_receipts(directory)
|
||||
return path
|
||||
except Exception as exc: # pragma: no cover - defensive
|
||||
@@ -227,194 +185,134 @@ def finalize_update_receipt(
|
||||
return None
|
||||
|
||||
|
||||
def finalize_pending_update_receipt(
|
||||
exit_code: Optional[int] = None, stop_reason: str = ""
|
||||
) -> Optional[Path]:
|
||||
"""Command-boundary safety net: persist a still-open receipt, if any.
|
||||
def finalize_pending_update_receipt(exit_code: Optional[int] = None, stop_reason: str = "") -> Optional[Path]:
|
||||
"""Command-boundary safety net: persist a still-open receipt, if any. Never raises.
|
||||
|
||||
``hermes update`` has many early ``sys.exit`` paths (preflight refusals, venv-holder
|
||||
refusal, fetch failure) predating the inner finalize calls; any receipt still open when the
|
||||
command unwinds is finalized here so refused/failed runs — where a receipt matters most —
|
||||
leave a record. No-op when nothing is open (inner paths finalize exactly-once via the popped
|
||||
singleton). Never raises. Exit 0/None → ``success``, exit 2 → ``refused`` (preflight
|
||||
convention), else → ``failed``.
|
||||
``hermes update`` has many early ``sys.exit`` paths (preflight refusals, venv-holder refusal,
|
||||
fetch failure) predating the inner finalize calls; finalizing here means refused/failed runs —
|
||||
where a receipt matters most — leave a record. Exit 0/None → ``success``, exit 2 → ``refused``
|
||||
(preflight convention), else → ``failed``.
|
||||
"""
|
||||
if _current is None:
|
||||
return None
|
||||
if exit_code in (0, None):
|
||||
outcome = "success"
|
||||
elif exit_code == 2:
|
||||
outcome = "refused"
|
||||
else:
|
||||
outcome = "failed"
|
||||
outcome = "success" if exit_code in (0, None) else "refused" if exit_code == 2 else "failed"
|
||||
if exit_code is not None:
|
||||
try:
|
||||
with suppress(Exception):
|
||||
_current.data["exit_code"] = int(exit_code)
|
||||
except Exception:
|
||||
pass
|
||||
return finalize_update_receipt(outcome, stop_reason=stop_reason)
|
||||
|
||||
|
||||
def _prune_old_receipts(directory: Path) -> None:
|
||||
try:
|
||||
receipts = sorted(
|
||||
(p for p in directory.glob("update_*.json") if p.is_file()),
|
||||
key=lambda p: p.stat().st_mtime,
|
||||
reverse=True,
|
||||
)
|
||||
for stale in receipts[_RECEIPT_KEEP:]:
|
||||
try:
|
||||
with suppress(Exception):
|
||||
receipts = (p for p in directory.glob("update_*.json") if p.is_file())
|
||||
for stale in sorted(receipts, key=lambda p: p.stat().st_mtime, reverse=True)[_RECEIPT_KEEP:]:
|
||||
with suppress(OSError):
|
||||
stale.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def read_latest_receipt() -> Optional[dict[str, Any]]:
|
||||
"""Read the most recent update receipt, or None. Never raises."""
|
||||
try:
|
||||
with suppress(Exception):
|
||||
path = _receipt_dir() / "latest.json"
|
||||
if not path.is_file():
|
||||
return None
|
||||
payload = json.loads(path.read_text(encoding="utf-8"))
|
||||
return payload if isinstance(payload, dict) else None
|
||||
except Exception:
|
||||
if path.is_file():
|
||||
payload = json.loads(path.read_text(encoding="utf-8"))
|
||||
return payload if isinstance(payload, dict) else None
|
||||
return None
|
||||
|
||||
|
||||
def _profile_homes() -> list[tuple[str, Path]]:
|
||||
"""``(profile, home)`` for the default home plus every valid named profile dir, sorted."""
|
||||
from hermes_cli.profiles import _get_default_hermes_home, _get_profiles_root, _PROFILE_ID_RE
|
||||
|
||||
homes: list[tuple[str, Path]] = []
|
||||
default_home = _get_default_hermes_home()
|
||||
if default_home.is_dir():
|
||||
homes.append(("default", default_home))
|
||||
root = _get_profiles_root()
|
||||
if root.is_dir():
|
||||
homes.extend(
|
||||
(entry.name, entry)
|
||||
for entry in sorted(root.iterdir())
|
||||
if entry.is_dir() and entry.name != "default" and _PROFILE_ID_RE.match(entry.name)
|
||||
)
|
||||
return homes
|
||||
|
||||
|
||||
def _socket_identity(home: Path) -> Optional[tuple[int, dict]]:
|
||||
"""``(pid, identity)`` declared by the gateway owning ``home``'s control socket, else None.
|
||||
|
||||
A live ``identify`` answer is authoritative — no PID-reuse or stale-file heuristics. Callers
|
||||
fall back to ``gateway_state.json`` for gateways that predate the socket or whose socket
|
||||
didn't bind.
|
||||
"""
|
||||
try:
|
||||
from gateway.control_socket import identify_gateway
|
||||
|
||||
identity = identify_gateway(home)
|
||||
return (int(identity.get("pid")), identity) if identity else None
|
||||
except Exception: # probe failure, no gateway, or an unparseable pid
|
||||
return None
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Fleet version verification
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def _sha_state(code_sha: Any, expected_sha: Any) -> str:
|
||||
if not code_sha or not expected_sha:
|
||||
return "unknown"
|
||||
return "current" if str(code_sha) == str(expected_sha) else "stale"
|
||||
|
||||
|
||||
def _fleet_row(profile: str, pid: int, code_sha: Any, code_version: Any, expected_sha: Any) -> dict[str, Any]:
|
||||
def _fleet_row(
|
||||
profile: str, pid: int, code_sha: Any, code_version: Any, expected_sha: Any, state: str = "unknown"
|
||||
) -> dict[str, Any]:
|
||||
if state == "unknown" and code_sha and expected_sha:
|
||||
state = "current" if str(code_sha) == str(expected_sha) else "stale"
|
||||
return {
|
||||
"profile": profile,
|
||||
"pid": pid,
|
||||
"code_sha": str(code_sha) if code_sha else None,
|
||||
"code_version": code_version,
|
||||
"state": _sha_state(code_sha, expected_sha),
|
||||
"profile": profile, "pid": pid, "code_sha": str(code_sha) if code_sha else None,
|
||||
"code_version": code_version, "state": state,
|
||||
}
|
||||
|
||||
|
||||
def collect_fleet_versions(
|
||||
*, pre_restart_pids: Optional[list[int]] = None
|
||||
) -> list[dict[str, Any]]:
|
||||
def collect_fleet_versions(*, pre_restart_pids: Optional[list[int]] = None) -> list[dict[str, Any]]:
|
||||
"""Snapshot every profile's gateway code identity vs. the current tree.
|
||||
|
||||
Rollout safety: ``down`` requires membership in ``pre_restart_pids`` — a stale state file from a
|
||||
long-dead gateway (machine reboot, manual kill weeks ago) must NOT fail every future update.
|
||||
Callers that don't have a pre-restart snapshot (``None``/empty) get the historical behavior:
|
||||
dead PIDs are skipped.
|
||||
Without a pre-restart snapshot (``None``/empty) dead PIDs are skipped (historical behavior).
|
||||
"""
|
||||
# Runtime-status states that mean "this record does not describe a
|
||||
# gateway that should be running now" — no down row for these.
|
||||
# Runtime-status states that do not describe a gateway that should be running now — no down row.
|
||||
_NOT_EXPECTED_STATES = {"stopped", "startup_failed"}
|
||||
_pre_restart = {int(p) for p in (pre_restart_pids or []) if isinstance(p, int)}
|
||||
results: list[dict[str, Any]] = []
|
||||
try:
|
||||
expected_sha = None
|
||||
with suppress(Exception):
|
||||
from hermes_cli.build_info import get_code_identity
|
||||
|
||||
expected_sha = (get_code_identity(refresh=True) or {}).get("sha")
|
||||
except Exception:
|
||||
expected_sha = None
|
||||
|
||||
try:
|
||||
from gateway.status import read_runtime_status, runtime_status_pid_is_live
|
||||
from hermes_cli.profiles import (
|
||||
_get_default_hermes_home,
|
||||
_get_profiles_root,
|
||||
_PROFILE_ID_RE,
|
||||
)
|
||||
|
||||
homes: list[tuple[str, Path]] = []
|
||||
default_home = _get_default_hermes_home()
|
||||
if default_home.is_dir():
|
||||
homes.append(("default", default_home))
|
||||
profiles_root = _get_profiles_root()
|
||||
if profiles_root.is_dir():
|
||||
for entry in sorted(profiles_root.iterdir()):
|
||||
if entry.is_dir() and entry.name != "default" and _PROFILE_ID_RE.match(entry.name):
|
||||
homes.append((entry.name, entry))
|
||||
|
||||
for profile, home in homes:
|
||||
# Prefer the gateway-owned control socket (#92091): a live
|
||||
# `identify` answer is authoritative — no PID-reuse or stale-file
|
||||
# heuristics. Fall back to gateway_state.json for gateways that
|
||||
# predate the socket or whose socket didn't bind.
|
||||
identity = None
|
||||
try:
|
||||
from gateway.control_socket import identify_gateway
|
||||
|
||||
identity = identify_gateway(home)
|
||||
except Exception:
|
||||
identity = None
|
||||
if identity:
|
||||
try:
|
||||
pid = int(identity.get("pid"))
|
||||
except (TypeError, ValueError):
|
||||
pid = None
|
||||
if pid is not None:
|
||||
row = _fleet_row(
|
||||
profile, pid, identity.get("code_sha"),
|
||||
identity.get("code_version"), expected_sha,
|
||||
)
|
||||
row["source"] = "socket"
|
||||
results.append(row)
|
||||
continue
|
||||
status_path = home / "gateway_state.json"
|
||||
record = read_runtime_status(status_path)
|
||||
for profile, home in _profile_homes():
|
||||
sock = _socket_identity(home)
|
||||
if sock is not None:
|
||||
pid, identity = sock
|
||||
row = _fleet_row(profile, pid, identity.get("code_sha"), identity.get("code_version"), expected_sha)
|
||||
results.append({**row, "source": "socket"})
|
||||
continue
|
||||
record = read_runtime_status(home / "gateway_state.json")
|
||||
if not record:
|
||||
continue
|
||||
pid = record.get("pid")
|
||||
try:
|
||||
pid = int(pid)
|
||||
pid = int(record.get("pid"))
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
if not runtime_status_pid_is_live(record):
|
||||
# Dead PID (or a live PID recycled by an unrelated process
|
||||
# during the update's own churn — #93258): a DOWN row only
|
||||
# when this exact pid was alive at update start AND the
|
||||
# record still claims a running state — "the restart phase
|
||||
# stopped it and nothing came back." Everything else (clean
|
||||
# stop, startup failure, stale record from a long-dead
|
||||
# gateway) keeps the historical no-row behavior so the
|
||||
# feature's rollout can't false-positive.
|
||||
#
|
||||
# ``_pre_restart`` is a bare set of PIDs, not (pid, start_time)
|
||||
# pairs, so a recycled PID from gateway A landing in B's stale
|
||||
# record could still mislabel B as down if A's PID happened to
|
||||
# be in the pre-restart snapshot — inherent to the snapshot's
|
||||
# data model, not something this guard can fix on its own.
|
||||
gw_state = record.get("gateway_state")
|
||||
if (
|
||||
pid in _pre_restart
|
||||
and isinstance(gw_state, str)
|
||||
and gw_state
|
||||
and gw_state not in _NOT_EXPECTED_STATES
|
||||
):
|
||||
results.append(
|
||||
{
|
||||
"profile": profile,
|
||||
"pid": pid,
|
||||
"code_sha": None,
|
||||
"code_version": record.get("code_version"),
|
||||
"state": "down",
|
||||
}
|
||||
)
|
||||
continue
|
||||
results.append(
|
||||
_fleet_row(
|
||||
profile, pid, record.get("code_sha"),
|
||||
record.get("code_version"), expected_sha,
|
||||
if runtime_status_pid_is_live(record):
|
||||
results.append(
|
||||
_fleet_row(profile, pid, record.get("code_sha"), record.get("code_version"), expected_sha)
|
||||
)
|
||||
)
|
||||
continue
|
||||
# Dead PID (or a live PID recycled by an unrelated process during the update's own
|
||||
# churn): a DOWN row only when this exact pid was alive at update start AND the record
|
||||
# still claims a running state — "the restart phase stopped it and nothing came back."
|
||||
# Everything else (clean stop, startup failure, long-dead stale record) keeps the no-row
|
||||
# behavior so the rollout can't false-positive. ``_pre_restart`` is a bare PID set, not
|
||||
# (pid, start_time) pairs, so a recycled PID from gateway A landing in B's stale record
|
||||
# could still mislabel B as down — inherent to the snapshot's data model.
|
||||
gw_state = record.get("gateway_state")
|
||||
if pid in _pre_restart and isinstance(gw_state, str) and gw_state and gw_state not in _NOT_EXPECTED_STATES:
|
||||
results.append(_fleet_row(profile, pid, None, record.get("code_version"), None, state="down"))
|
||||
except Exception as exc:
|
||||
logger.debug("Fleet version probe failed: %s", exc)
|
||||
return results
|
||||
@@ -426,21 +324,18 @@ def print_fleet_version_matrix(fleet: list[dict[str, Any]]) -> bool:
|
||||
Returns True when at least one gateway is provably stale (still serving pre-update code) OR
|
||||
provably down (killed by the restart phase, nothing came back), so the caller can escalate.
|
||||
``unknown`` entries are reported but do NOT fail the update: gateways started before the
|
||||
code-identity stamp existed have no sha to compare, and failing them would be a false-
|
||||
positive storm.
|
||||
code-identity stamp existed have no sha to compare, and failing them would be a false-positive
|
||||
storm.
|
||||
"""
|
||||
if not fleet:
|
||||
return False
|
||||
any_stale = False
|
||||
any_down = False
|
||||
any_stale = any_down = False
|
||||
print()
|
||||
print("Fleet version check:")
|
||||
for entry in fleet:
|
||||
sha = entry.get("code_sha")
|
||||
short = sha[:8] if isinstance(sha, str) and sha else "?"
|
||||
state = entry.get("state")
|
||||
profile = entry.get("profile")
|
||||
pid = entry.get("pid")
|
||||
state, profile, pid = entry.get("state"), entry.get("profile"), entry.get("pid")
|
||||
if state == "current":
|
||||
print(f" ✓ {profile} (pid {pid}) @ {short} — up to date")
|
||||
elif state == "stale":
|
||||
|
||||
Reference in New Issue
Block a user