fix(update): conservative outcomes + serve-ledger coverage for fresh restart recovery
Salvage adjustments to PR #94392 per review: - Narrow the supervisor claim to the systemd-VERIFIED path only. The fresh recovery child now probes 'systemctl --user is-active' after each relaunch; only an observed-active systemd unit is reported 'verified'. A relaunch that merely exited 0 is labelled 'relaunch_attempted', never counts as supervisor coverage, and never clears gateway_fleet_restart_incomplete. - Serve-owned runtimes (serve/dashboard entries from the spawn ledger, per the update_inventory serve collector) are no longer silently skipped: the recovery pass records them (and manual gateways) as skipped-with-reason in the recovery result and the persisted update receipt. - Receipt fresh_recovery persists the conservative vocabulary (requested/verified/relaunch_attempted/failed/skipped); 'succeeded' is gone. - Added an end-to-end test that drives the real recovery module in a genuinely fresh interpreter (sitecustomize shim intercepts the grandchild 'gateway restart' and systemctl probes).
This commit is contained in:
@@ -5083,6 +5083,7 @@ _LAZY_COMMAND_EXPORTS = {
|
||||
"_format_concurrent_instances_message",
|
||||
"_format_time_ago",
|
||||
"_gateway_service_matches_profile",
|
||||
"_gateway_recovery_partition",
|
||||
"_gateway_restart_recovery_profiles",
|
||||
"_handoff_reapable_backend_pids",
|
||||
"_ledger_reapable_backend_pids",
|
||||
|
||||
+143
-58
@@ -5615,36 +5615,89 @@ def _gateway_service_matches_profile(profile: str, service: object) -> bool:
|
||||
}
|
||||
|
||||
|
||||
def _gateway_restart_recovery_profiles(
|
||||
def _gateway_recovery_partition(
|
||||
plan, *, skip_profiles: set[str] | None = None
|
||||
) -> list[str]:
|
||||
"""Return supervised gateway profiles that a fresh process may restart.
|
||||
) -> tuple[dict[str, str], list[dict]]:
|
||||
"""Partition pre-update runtimes into fresh-restart candidates and skips.
|
||||
|
||||
The update inventory is captured before the checkout changes. It is the
|
||||
only safe source here: re-importing ``hermes_cli.gateway`` in the failing
|
||||
interpreter is exactly what can raise the original ``ImportError``. Manual
|
||||
gateways are intentionally excluded because no supervisor is available to
|
||||
bring them back after a stop.
|
||||
interpreter is exactly what can raise the original ``ImportError``.
|
||||
|
||||
Returns ``(candidates, skipped)`` where ``candidates`` maps profile →
|
||||
supervisor for supervised gateway runtimes the fresh process may restart,
|
||||
and ``skipped`` lists every other inventoried runtime the recovery pass
|
||||
deliberately does NOT touch, each with an explicit reason. Nothing from
|
||||
the spawn ledger may vanish from the recovery pass silently: manual
|
||||
gateways have no relaunch authority, and serve/dashboard runtimes (the
|
||||
``update_inventory`` serve collector) are owned by the Desktop app or a
|
||||
human terminal, not by this recovery boundary.
|
||||
"""
|
||||
profiles: set[str] = set()
|
||||
skip_profiles = skip_profiles or set()
|
||||
candidates: dict[str, str] = {}
|
||||
skipped: list[dict] = []
|
||||
try:
|
||||
for runtime in getattr(plan, "runtimes", ()) or ():
|
||||
if (
|
||||
getattr(runtime, "kind", None) == "gateway"
|
||||
and getattr(runtime, "supervisor", None) in _FRESH_RESTART_SUPERVISORS
|
||||
):
|
||||
profile = getattr(runtime, "profile", None)
|
||||
if isinstance(profile, str) and profile and profile not in skip_profiles:
|
||||
profiles.add(profile)
|
||||
kind = getattr(runtime, "kind", None)
|
||||
profile = getattr(runtime, "profile", None)
|
||||
supervisor = getattr(runtime, "supervisor", None)
|
||||
if not isinstance(profile, str) or not profile:
|
||||
continue
|
||||
if kind == "gateway":
|
||||
if profile in skip_profiles:
|
||||
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"
|
||||
),
|
||||
}
|
||||
)
|
||||
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:
|
||||
reason = (
|
||||
"manually launched serve/dashboard has no relaunch"
|
||||
" authority; left running for explicit operator"
|
||||
" restart"
|
||||
)
|
||||
skipped.append(
|
||||
{
|
||||
"profile": profile,
|
||||
"kind": str(kind),
|
||||
"supervisor": str(supervisor),
|
||||
"reason": reason,
|
||||
}
|
||||
)
|
||||
except Exception as exc:
|
||||
logger.debug("Could not prepare fresh gateway restart profiles: %s", exc)
|
||||
return sorted(profiles)
|
||||
return candidates, skipped
|
||||
|
||||
|
||||
def _gateway_restart_recovery_profiles(
|
||||
plan, *, skip_profiles: set[str] | None = None
|
||||
) -> list[str]:
|
||||
"""Return supervised gateway profiles that a fresh process may restart."""
|
||||
candidates, _ = _gateway_recovery_partition(plan, skip_profiles=skip_profiles)
|
||||
return sorted(candidates)
|
||||
|
||||
|
||||
def _recover_gateway_restart_after_abort(
|
||||
plan, *, gateway_mode: bool, skip_profiles: set[str] | None = None
|
||||
) -> dict[str, list[str]]:
|
||||
) -> dict[str, list]:
|
||||
"""Retry supervised gateway restarts from a clean Python process.
|
||||
|
||||
``hermes update`` normally performs the fleet restart in the interpreter
|
||||
@@ -5658,12 +5711,38 @@ def _recover_gateway_restart_after_abort(
|
||||
Only profiles classified as supervisor-owned by the pre-update inventory
|
||||
are handed off. A manual gateway must remain running and be reported for
|
||||
explicit operator action rather than being killed without a relaunch
|
||||
authority. The returned protocol is persisted in the update receipt so
|
||||
operators can distinguish a spawn failure from a per-profile failure.
|
||||
authority; serve/dashboard runtimes from the spawn ledger are likewise
|
||||
recorded as skipped with a reason instead of vanishing from the pass.
|
||||
The returned protocol is persisted in the update receipt so operators can
|
||||
distinguish a spawn failure from a per-profile failure.
|
||||
|
||||
Outcome honesty: ``verified`` means the fresh child independently observed
|
||||
the profile's systemd unit active after the relaunch. A zero exit from
|
||||
``gateway restart`` alone is NOT observed proof that the new code
|
||||
generation is serving, so those outcomes are reported as
|
||||
``relaunch_attempted`` and never claim supervisor coverage.
|
||||
"""
|
||||
profiles = _gateway_restart_recovery_profiles(plan, skip_profiles=skip_profiles)
|
||||
candidates, skipped = _gateway_recovery_partition(
|
||||
plan, skip_profiles=skip_profiles
|
||||
)
|
||||
profiles = sorted(candidates)
|
||||
if not profiles:
|
||||
return {"requested": [], "succeeded": [], "failed": []}
|
||||
return {
|
||||
"requested": [],
|
||||
"verified": [],
|
||||
"relaunch_attempted": [],
|
||||
"failed": [],
|
||||
"skipped": skipped,
|
||||
}
|
||||
|
||||
def _all_failed() -> dict[str, list]:
|
||||
return {
|
||||
"requested": profiles,
|
||||
"verified": [],
|
||||
"relaunch_attempted": [],
|
||||
"failed": profiles,
|
||||
"skipped": skipped,
|
||||
}
|
||||
|
||||
command = [
|
||||
sys.executable,
|
||||
@@ -5685,7 +5764,7 @@ def _recover_gateway_restart_after_abort(
|
||||
systemd_run = shutil.which("systemd-run")
|
||||
if not systemd_run:
|
||||
logger.warning("Cannot isolate fresh gateway recovery from the gateway cgroup")
|
||||
return {"requested": profiles, "succeeded": [], "failed": profiles}
|
||||
return _all_failed()
|
||||
command = [
|
||||
systemd_run,
|
||||
"--user",
|
||||
@@ -5697,7 +5776,7 @@ def _recover_gateway_restart_after_abort(
|
||||
]
|
||||
|
||||
kwargs = {
|
||||
"input": json.dumps({"profiles": profiles}),
|
||||
"input": json.dumps({"profiles": profiles, "supervisors": candidates}),
|
||||
"capture_output": True,
|
||||
"text": True,
|
||||
"encoding": "utf-8",
|
||||
@@ -5718,39 +5797,52 @@ def _recover_gateway_restart_after_abort(
|
||||
result = subprocess.run(command, **kwargs)
|
||||
except (OSError, subprocess.TimeoutExpired) as exc:
|
||||
logger.warning("Fresh gateway restart recovery failed: %s", exc)
|
||||
return {"requested": profiles, "succeeded": [], "failed": profiles}
|
||||
return _all_failed()
|
||||
|
||||
if result.returncode != 0:
|
||||
logger.warning("Fresh gateway restart recovery exited %s", result.returncode)
|
||||
return {"requested": profiles, "succeeded": [], "failed": profiles}
|
||||
return _all_failed()
|
||||
|
||||
try:
|
||||
recovery_result = json.loads(result.stdout or "")
|
||||
succeeded = recovery_result.get("succeeded")
|
||||
verified = recovery_result.get("verified")
|
||||
relaunch_attempted = recovery_result.get("relaunch_attempted")
|
||||
failed = recovery_result.get("failed")
|
||||
except (AttributeError, TypeError, ValueError):
|
||||
logger.warning("Fresh gateway restart recovery returned invalid JSON")
|
||||
return {"requested": profiles, "succeeded": [], "failed": profiles}
|
||||
return _all_failed()
|
||||
|
||||
buckets = (verified, relaunch_attempted, failed)
|
||||
reported: list[str] = []
|
||||
if all(isinstance(bucket, list) for bucket in buckets):
|
||||
reported = [*verified, *relaunch_attempted, *failed]
|
||||
if (
|
||||
not isinstance(succeeded, list)
|
||||
or not isinstance(failed, list)
|
||||
or any(not isinstance(profile, str) for profile in [*succeeded, *failed])
|
||||
or set(succeeded) != set(profiles)
|
||||
or failed
|
||||
not all(isinstance(bucket, list) for bucket in buckets)
|
||||
or any(not isinstance(profile, str) for profile in reported)
|
||||
or set(reported) != set(profiles)
|
||||
or len(reported) != len(set(reported))
|
||||
):
|
||||
logger.warning("Fresh gateway restart recovery returned incomplete profiles")
|
||||
return {
|
||||
"requested": profiles,
|
||||
"succeeded": succeeded if isinstance(succeeded, list) else [],
|
||||
"failed": failed if isinstance(failed, list) else profiles,
|
||||
}
|
||||
return _all_failed()
|
||||
|
||||
print(
|
||||
" ✓ Retried supervised gateway restart(s) in a fresh process: "
|
||||
+ ", ".join(profiles)
|
||||
)
|
||||
return {"requested": profiles, "succeeded": succeeded, "failed": []}
|
||||
if verified:
|
||||
print(
|
||||
" ✓ Restarted supervised gateway(s) in a fresh process"
|
||||
" (systemd-verified active): " + ", ".join(sorted(verified))
|
||||
)
|
||||
if relaunch_attempted:
|
||||
print(
|
||||
" ⚠ Relaunch attempted in a fresh process but not"
|
||||
" supervisor-verified (check these gateways manually): "
|
||||
+ ", ".join(sorted(relaunch_attempted))
|
||||
)
|
||||
return {
|
||||
"requested": profiles,
|
||||
"verified": sorted(verified),
|
||||
"relaunch_attempted": sorted(relaunch_attempted),
|
||||
"failed": sorted(failed),
|
||||
"skipped": skipped,
|
||||
}
|
||||
|
||||
|
||||
def _warn_gateway_restart_phase_aborted(exc: BaseException, pids) -> None:
|
||||
@@ -8802,16 +8894,14 @@ def _cmd_update_impl(args, gateway_mode: bool):
|
||||
gateway_mode=gateway_mode,
|
||||
skip_profiles=_already_restarted_profiles,
|
||||
)
|
||||
_recovery_profiles = list(_recovery_result.get("requested") or [])
|
||||
_recovery_succeeded = set(_recovery_result.get("succeeded") or [])
|
||||
_recovered = not _recovery_profiles or (
|
||||
_recovery_succeeded == set(_recovery_profiles)
|
||||
and not _recovery_result.get("failed")
|
||||
)
|
||||
if _recovery_succeeded:
|
||||
# Only systemd-VERIFIED outcomes may claim supervisor coverage.
|
||||
# A relaunch that merely exited 0 ("relaunch_attempted") was never
|
||||
# observed by the code and must not clear the incomplete flag.
|
||||
_recovery_verified = set(_recovery_result.get("verified") or [])
|
||||
if _recovery_verified:
|
||||
relaunched_profiles.extend(
|
||||
profile
|
||||
for profile in sorted(_recovery_succeeded)
|
||||
for profile in sorted(_recovery_verified)
|
||||
if profile not in relaunched_profiles
|
||||
)
|
||||
_planned_gateway_runtimes = [
|
||||
@@ -8823,18 +8913,13 @@ def _cmd_update_impl(args, gateway_mode: bool):
|
||||
_planned_gateway_profiles = {
|
||||
runtime.profile for runtime in _planned_gateway_runtimes
|
||||
}
|
||||
_covered_gateway_profiles = _already_restarted_profiles | set(
|
||||
_recovery_succeeded
|
||||
)
|
||||
_all_gateway_profiles_are_safe_to_recover = all(
|
||||
runtime.profile in _already_restarted_profiles
|
||||
or getattr(runtime, "supervisor", None) in _FRESH_RESTART_SUPERVISORS
|
||||
for runtime in _planned_gateway_runtimes
|
||||
_covered_gateway_profiles = (
|
||||
_already_restarted_profiles | _recovery_verified
|
||||
)
|
||||
_recovery_complete = bool(_planned_gateway_profiles) and (
|
||||
_planned_gateway_profiles <= _covered_gateway_profiles
|
||||
and _all_gateway_profiles_are_safe_to_recover
|
||||
and (not _recovery_profiles or _recovered)
|
||||
and not _recovery_result.get("failed")
|
||||
and not _recovery_result.get("relaunch_attempted")
|
||||
)
|
||||
if _recovery_complete:
|
||||
# The fresh child is the recovery terminal result. Leave the
|
||||
|
||||
@@ -113,17 +113,26 @@ class UpdateReceipt:
|
||||
"phase_error": phase_error,
|
||||
}
|
||||
if fresh_recovery is not None:
|
||||
result["fresh_recovery"] = {
|
||||
"requested": [
|
||||
str(profile) for profile in fresh_recovery.get("requested", [])
|
||||
],
|
||||
"succeeded": [
|
||||
str(profile) for profile in fresh_recovery.get("succeeded", [])
|
||||
],
|
||||
"failed": [
|
||||
str(profile) for profile in fresh_recovery.get("failed", [])
|
||||
],
|
||||
# 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"] = [
|
||||
{
|
||||
"profile": str(entry.get("profile", "")),
|
||||
"kind": str(entry.get("kind", "")),
|
||||
"supervisor": str(entry.get("supervisor", "")),
|
||||
"reason": str(entry.get("reason", "")),
|
||||
}
|
||||
for entry in fresh_recovery.get("skipped", [])
|
||||
if isinstance(entry, dict)
|
||||
]
|
||||
result["fresh_recovery"] = persisted
|
||||
self.data["gateway_restart"] = result
|
||||
|
||||
def finalize(self, outcome: str) -> None:
|
||||
|
||||
@@ -6,6 +6,19 @@ itself and launches the regular per-profile gateway command in a new
|
||||
interpreter. It is used only after the in-process restart phase has raised, so
|
||||
that the recovery path cannot inherit the stale ``sys.modules`` graph that
|
||||
caused the failure.
|
||||
|
||||
Outcome vocabulary (deliberately conservative):
|
||||
|
||||
- ``verified`` — the relaunch command exited 0 AND the profile's
|
||||
systemd unit was independently observed ``active`` afterwards. This is the
|
||||
only outcome that may claim supervisor coverage.
|
||||
- ``relaunch_attempted`` — the relaunch command exited 0 but no independent
|
||||
supervisor observation was possible (non-systemd supervisor, ``systemctl``
|
||||
missing, or the unit probe was inconclusive). ``rc == 0`` from
|
||||
``gateway restart`` is not proof that the new code generation is running,
|
||||
so this outcome must never be treated as verified coverage.
|
||||
- ``failed`` — the relaunch command errored, timed out, or exited
|
||||
non-zero.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -14,15 +27,18 @@ import argparse
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
from collections.abc import Callable, Iterable
|
||||
from collections.abc import Callable, Iterable, Mapping
|
||||
from typing import Any
|
||||
|
||||
_RECOVERY_ENV = "HERMES_UPDATE_RESTART_RECOVERY"
|
||||
_GATEWAY_MARKERS = ("_HERMES_GATEWAY", "HERMES_GATEWAY", "HERMES_GATEWAY_MODE")
|
||||
_PROFILE_RESTART_TIMEOUT = 90
|
||||
_VERIFY_TIMEOUT = 15
|
||||
_PROFILE_ID_RE = re.compile(r"^[a-z0-9][a-z0-9_-]{0,63}$")
|
||||
_SUPERVISOR_RE = re.compile(r"^[a-z0-9][a-z0-9_-]{0,31}$")
|
||||
|
||||
|
||||
def _profile_command(profile: str) -> list[str]:
|
||||
@@ -78,8 +94,57 @@ def _run_profile_restart(
|
||||
return getattr(result, "returncode", 1) == 0
|
||||
|
||||
|
||||
def _systemd_unit_candidates(profile: str) -> tuple[str, ...]:
|
||||
"""Unit names the existing systemd gateway lifecycle produces per profile."""
|
||||
if profile == "default":
|
||||
return (
|
||||
"hermes-gateway.service",
|
||||
"gateway.service",
|
||||
"gateway-default.service",
|
||||
)
|
||||
return (
|
||||
f"hermes-gateway-{profile}.service",
|
||||
f"gateway-{profile}.service",
|
||||
)
|
||||
|
||||
|
||||
def _systemd_verified_active(profile: str, *, run: Callable[..., Any]) -> bool:
|
||||
"""Return True only when systemd itself reports the profile's unit active.
|
||||
|
||||
This is the observation that separates ``verified`` from
|
||||
``relaunch_attempted``. Any failure here (no ``systemctl``, probe error,
|
||||
unit not ``active``) means we could NOT verify — never that the restart
|
||||
failed.
|
||||
"""
|
||||
systemctl = shutil.which("systemctl")
|
||||
if not systemctl:
|
||||
return False
|
||||
for unit in _systemd_unit_candidates(profile):
|
||||
try:
|
||||
result = run(
|
||||
[systemctl, "--user", "is-active", unit],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
encoding="utf-8",
|
||||
errors="replace",
|
||||
check=False,
|
||||
timeout=_VERIFY_TIMEOUT,
|
||||
)
|
||||
except (OSError, subprocess.TimeoutExpired):
|
||||
continue
|
||||
if (
|
||||
getattr(result, "returncode", 1) == 0
|
||||
and (getattr(result, "stdout", "") or "").strip() == "active"
|
||||
):
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
def restart_profiles(
|
||||
profiles: Iterable[str], *, run: Callable[..., Any] = subprocess.run
|
||||
profiles: Iterable[str],
|
||||
*,
|
||||
supervisors: Mapping[str, str] | None = None,
|
||||
run: Callable[..., Any] = subprocess.run,
|
||||
) -> dict[str, list[str]]:
|
||||
"""Restart the supplied profiles and return per-profile terminal results.
|
||||
|
||||
@@ -87,21 +152,38 @@ def restart_profiles(
|
||||
supervisor. Manual gateways are intentionally excluded before this module
|
||||
is called: killing one without a relaunch authority would turn stale code
|
||||
into an outage.
|
||||
|
||||
A profile only lands in ``verified`` when its supervisor is systemd and
|
||||
``systemctl --user is-active`` independently confirms the unit after the
|
||||
relaunch command succeeded. Every other zero-exit relaunch is reported as
|
||||
``relaunch_attempted`` — the code cannot observe supervisor coverage for
|
||||
those paths and must not claim it.
|
||||
"""
|
||||
supervisors = supervisors or {}
|
||||
normalized = sorted(
|
||||
{profile for profile in profiles if isinstance(profile, str) and profile}
|
||||
)
|
||||
succeeded: list[str] = []
|
||||
verified: list[str] = []
|
||||
relaunch_attempted: list[str] = []
|
||||
failed: list[str] = []
|
||||
for profile in normalized:
|
||||
if _run_profile_restart(profile, run=run):
|
||||
succeeded.append(profile)
|
||||
else:
|
||||
if not _run_profile_restart(profile, run=run):
|
||||
failed.append(profile)
|
||||
return {"succeeded": succeeded, "failed": failed}
|
||||
continue
|
||||
if supervisors.get(profile) == "systemd" and _systemd_verified_active(
|
||||
profile, run=run
|
||||
):
|
||||
verified.append(profile)
|
||||
else:
|
||||
relaunch_attempted.append(profile)
|
||||
return {
|
||||
"verified": verified,
|
||||
"relaunch_attempted": relaunch_attempted,
|
||||
"failed": failed,
|
||||
}
|
||||
|
||||
|
||||
def _parse_payload(stream) -> list[str]:
|
||||
def _parse_payload(stream) -> tuple[list[str], dict[str, str]]:
|
||||
payload = json.load(stream)
|
||||
profiles = payload.get("profiles") if isinstance(payload, dict) else None
|
||||
if not isinstance(profiles, list):
|
||||
@@ -111,7 +193,19 @@ def _parse_payload(stream) -> list[str]:
|
||||
for profile in profiles
|
||||
):
|
||||
raise ValueError("recovery profiles contain an invalid profile id")
|
||||
return profiles
|
||||
raw_supervisors = payload.get("supervisors") if isinstance(payload, dict) else None
|
||||
supervisors: dict[str, str] = {}
|
||||
if raw_supervisors is not None:
|
||||
if not isinstance(raw_supervisors, dict) or any(
|
||||
not isinstance(profile, str)
|
||||
or not isinstance(supervisor, str)
|
||||
or not _PROFILE_ID_RE.fullmatch(profile)
|
||||
or not _SUPERVISOR_RE.fullmatch(supervisor)
|
||||
for profile, supervisor in raw_supervisors.items()
|
||||
):
|
||||
raise ValueError("recovery supervisors map is invalid")
|
||||
supervisors = dict(raw_supervisors)
|
||||
return profiles, supervisors
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
@@ -126,10 +220,19 @@ def main(argv: list[str] | None = None) -> int:
|
||||
parser.error("this command is an internal update-recovery entry point")
|
||||
|
||||
try:
|
||||
profiles = _parse_payload(sys.stdin)
|
||||
result = restart_profiles(profiles)
|
||||
profiles, supervisors = _parse_payload(sys.stdin)
|
||||
result = restart_profiles(profiles, supervisors=supervisors)
|
||||
except (ValueError, json.JSONDecodeError) as exc:
|
||||
print(json.dumps({"error": str(exc), "succeeded": [], "failed": []}))
|
||||
print(
|
||||
json.dumps(
|
||||
{
|
||||
"error": str(exc),
|
||||
"verified": [],
|
||||
"relaunch_attempted": [],
|
||||
"failed": [],
|
||||
}
|
||||
)
|
||||
)
|
||||
return 2
|
||||
|
||||
print(json.dumps(result, sort_keys=True))
|
||||
|
||||
@@ -68,9 +68,18 @@ class TestReceiptLifecycle:
|
||||
|
||||
def test_fresh_recovery_result_reaches_persisted_receipt(self, receipt_home):
|
||||
recovery = {
|
||||
"requested": ["coder", "default"],
|
||||
"succeeded": ["default"],
|
||||
"requested": ["coder", "default", "ops"],
|
||||
"verified": ["default"],
|
||||
"relaunch_attempted": ["ops"],
|
||||
"failed": ["coder"],
|
||||
"skipped": [
|
||||
{
|
||||
"profile": "desk",
|
||||
"kind": "serve",
|
||||
"supervisor": "desktop",
|
||||
"reason": "desktop app owns and respawns this serve backend",
|
||||
}
|
||||
],
|
||||
}
|
||||
|
||||
ur.begin_update_receipt()
|
||||
@@ -83,7 +92,11 @@ class TestReceiptLifecycle:
|
||||
path = _finalize("partial")
|
||||
|
||||
payload = json.loads(path.read_text(encoding="utf-8"))
|
||||
assert payload["gateway_restart"]["fresh_recovery"] == recovery
|
||||
persisted = payload["gateway_restart"]["fresh_recovery"]
|
||||
assert persisted == recovery
|
||||
# The conservative vocabulary is the persisted contract: no bucket may
|
||||
# rebrand an unverified relaunch as supervisor-backed success.
|
||||
assert "succeeded" not in persisted
|
||||
|
||||
def test_latest_pointer_written_and_readable(self, receipt_home):
|
||||
ur.begin_update_receipt()
|
||||
|
||||
@@ -3,8 +3,10 @@
|
||||
The updater may have loaded the pre-pull module graph when the checkout changes.
|
||||
If the in-process gateway restart phase then raises, retrying through the same
|
||||
interpreter cannot establish a coherent module generation. Recovery must use a
|
||||
new interpreter and must not invent a restart for manual gateways that have no
|
||||
supervisor to bring them back.
|
||||
new interpreter, must not invent a restart for manual gateways that have no
|
||||
supervisor to bring them back, and must not claim supervisor coverage it never
|
||||
observed: only a systemd-verified unit counts as ``verified``; a bare rc==0
|
||||
relaunch is ``relaunch_attempted``.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -14,6 +16,7 @@ import io
|
||||
import json
|
||||
import subprocess
|
||||
import sys
|
||||
import textwrap
|
||||
from types import SimpleNamespace
|
||||
|
||||
from hermes_cli import update_cmd
|
||||
@@ -26,10 +29,19 @@ class _Completed:
|
||||
self.stderr = stderr
|
||||
|
||||
|
||||
def _successful_recovery_result(profiles: list[str]) -> _Completed:
|
||||
def _successful_recovery_result(
|
||||
verified: list[str] | None = None,
|
||||
relaunch_attempted: list[str] | None = None,
|
||||
) -> _Completed:
|
||||
return _Completed(
|
||||
0,
|
||||
stdout=json.dumps({"succeeded": profiles, "failed": []}),
|
||||
stdout=json.dumps(
|
||||
{
|
||||
"verified": verified or [],
|
||||
"relaunch_attempted": relaunch_attempted or [],
|
||||
"failed": [],
|
||||
}
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
@@ -47,7 +59,7 @@ def test_abort_recovery_hands_managed_profiles_to_a_fresh_process(monkeypatch):
|
||||
|
||||
def fake_run(argv, **kwargs):
|
||||
calls.append((argv, kwargs))
|
||||
return _successful_recovery_result(["coder", "default"])
|
||||
return _successful_recovery_result(verified=["coder", "default"])
|
||||
|
||||
monkeypatch.setattr(update_cmd.subprocess, "run", fake_run)
|
||||
plan = SimpleNamespace(
|
||||
@@ -60,18 +72,24 @@ def test_abort_recovery_hands_managed_profiles_to_a_fresh_process(monkeypatch):
|
||||
)
|
||||
|
||||
result = update_cmd._recover_gateway_restart_after_abort(plan, gateway_mode=False)
|
||||
assert result == {
|
||||
"requested": ["coder", "default"],
|
||||
"succeeded": ["coder", "default"],
|
||||
"failed": [],
|
||||
}
|
||||
assert result["requested"] == ["coder", "default"]
|
||||
assert result["verified"] == ["coder", "default"]
|
||||
assert result["relaunch_attempted"] == []
|
||||
assert result["failed"] == []
|
||||
# Runtimes the pass does not own are recorded, not silently dropped.
|
||||
skipped = {(entry["profile"], entry["kind"]) for entry in result["skipped"]}
|
||||
assert skipped == {("manual-box", "gateway"), ("desktop", "serve")}
|
||||
assert all(entry["reason"] for entry in result["skipped"])
|
||||
|
||||
assert len(calls) == 1
|
||||
argv, kwargs = calls[0]
|
||||
assert argv[0] == sys.executable
|
||||
assert argv[1:4] == ["-m", "hermes_cli.update_restart_recovery", "--stdin"]
|
||||
payload = json.loads(kwargs["input"])
|
||||
assert payload == {"profiles": ["coder", "default"]}
|
||||
assert payload == {
|
||||
"profiles": ["coder", "default"],
|
||||
"supervisors": {"coder": "launchd", "default": "systemd"},
|
||||
}
|
||||
assert kwargs["text"] is True
|
||||
assert kwargs["capture_output"] is True
|
||||
assert kwargs["check"] is False
|
||||
@@ -88,16 +106,34 @@ def test_abort_recovery_does_not_claim_success_when_fresh_process_fails(monkeypa
|
||||
|
||||
result = update_cmd._recover_gateway_restart_after_abort(plan, gateway_mode=False)
|
||||
assert result["requested"] == ["default"]
|
||||
assert result["succeeded"] == []
|
||||
assert result["verified"] == []
|
||||
assert result["relaunch_attempted"] == []
|
||||
assert result["failed"] == ["default"]
|
||||
|
||||
|
||||
def test_abort_recovery_reports_unverified_relaunch_conservatively(monkeypatch):
|
||||
"""rc==0 without a systemd observation must not be reported as verified."""
|
||||
monkeypatch.setattr(
|
||||
update_cmd.subprocess,
|
||||
"run",
|
||||
lambda *args, **kwargs: _successful_recovery_result(
|
||||
relaunch_attempted=["default"]
|
||||
),
|
||||
)
|
||||
plan = SimpleNamespace(runtimes=[_runtime("default", "systemd")])
|
||||
|
||||
result = update_cmd._recover_gateway_restart_after_abort(plan, gateway_mode=False)
|
||||
assert result["verified"] == []
|
||||
assert result["relaunch_attempted"] == ["default"]
|
||||
assert result["failed"] == []
|
||||
|
||||
|
||||
def test_abort_recovery_skips_profiles_already_restarted_by_the_phase(monkeypatch):
|
||||
calls = []
|
||||
|
||||
def fake_run(argv, **kwargs):
|
||||
calls.append((argv, kwargs))
|
||||
return _successful_recovery_result(["coder"])
|
||||
return _successful_recovery_result(verified=["coder"])
|
||||
|
||||
monkeypatch.setattr(update_cmd.subprocess, "run", fake_run)
|
||||
plan = SimpleNamespace(
|
||||
@@ -109,8 +145,10 @@ def test_abort_recovery_skips_profiles_already_restarted_by_the_phase(monkeypatc
|
||||
gateway_mode=False,
|
||||
skip_profiles={"default"},
|
||||
)
|
||||
assert result == {"requested": ["coder"], "succeeded": ["coder"], "failed": []}
|
||||
assert json.loads(calls[0][1]["input"]) == {"profiles": ["coder"]}
|
||||
assert result["requested"] == ["coder"]
|
||||
assert result["verified"] == ["coder"]
|
||||
assert result["failed"] == []
|
||||
assert json.loads(calls[0][1]["input"])["profiles"] == ["coder"]
|
||||
|
||||
|
||||
def test_abort_recovery_rejects_partial_json_success(monkeypatch):
|
||||
@@ -119,7 +157,9 @@ def test_abort_recovery_rejects_partial_json_success(monkeypatch):
|
||||
"run",
|
||||
lambda *args, **kwargs: _Completed(
|
||||
0,
|
||||
stdout=json.dumps({"succeeded": ["default"], "failed": ["coder"]}),
|
||||
stdout=json.dumps(
|
||||
{"verified": ["default"], "relaunch_attempted": [], "failed": []}
|
||||
),
|
||||
),
|
||||
)
|
||||
plan = SimpleNamespace(
|
||||
@@ -127,9 +167,10 @@ def test_abort_recovery_rejects_partial_json_success(monkeypatch):
|
||||
)
|
||||
|
||||
result = update_cmd._recover_gateway_restart_after_abort(plan, gateway_mode=False)
|
||||
# "coder" is unaccounted for in the child's report: fail closed.
|
||||
assert result["requested"] == ["coder", "default"]
|
||||
assert result["succeeded"] == ["default"]
|
||||
assert result["failed"] == ["coder"]
|
||||
assert result["verified"] == []
|
||||
assert result["failed"] == ["coder", "default"]
|
||||
|
||||
|
||||
def test_abort_recovery_rejects_malformed_json_success(monkeypatch):
|
||||
@@ -142,7 +183,7 @@ def test_abort_recovery_rejects_malformed_json_success(monkeypatch):
|
||||
|
||||
result = update_cmd._recover_gateway_restart_after_abort(plan, gateway_mode=False)
|
||||
assert result["requested"] == ["default"]
|
||||
assert result["succeeded"] == []
|
||||
assert result["verified"] == []
|
||||
assert result["failed"] == ["default"]
|
||||
|
||||
|
||||
@@ -156,10 +197,36 @@ def test_abort_recovery_does_not_restart_manual_only_fleet(monkeypatch):
|
||||
plan = SimpleNamespace(runtimes=[_runtime("manual-box", "manual")])
|
||||
|
||||
result = update_cmd._recover_gateway_restart_after_abort(plan, gateway_mode=False)
|
||||
assert result == {"requested": [], "succeeded": [], "failed": []}
|
||||
assert result["requested"] == []
|
||||
assert result["verified"] == []
|
||||
assert result["failed"] == []
|
||||
assert [entry["profile"] for entry in result["skipped"]] == ["manual-box"]
|
||||
assert calls == []
|
||||
|
||||
|
||||
def test_abort_recovery_records_serve_runtimes_as_skipped_with_reason(monkeypatch):
|
||||
"""Serve/dashboard ledger entries must not vanish from the recovery pass."""
|
||||
monkeypatch.setattr(
|
||||
update_cmd.subprocess,
|
||||
"run",
|
||||
lambda *args, **kwargs: _successful_recovery_result(verified=["default"]),
|
||||
)
|
||||
plan = SimpleNamespace(
|
||||
runtimes=[
|
||||
_runtime("default", "systemd"),
|
||||
_runtime("default", "desktop", kind="serve"),
|
||||
_runtime("ops", "manual-serve", kind="dashboard"),
|
||||
]
|
||||
)
|
||||
|
||||
result = update_cmd._recover_gateway_restart_after_abort(plan, gateway_mode=False)
|
||||
by_kind = {entry["kind"]: entry for entry in result["skipped"]}
|
||||
assert set(by_kind) == {"serve", "dashboard"}
|
||||
assert "desktop app" in by_kind["serve"]["reason"]
|
||||
assert by_kind["dashboard"]["profile"] == "ops"
|
||||
assert "relaunch authority" in by_kind["dashboard"]["reason"]
|
||||
|
||||
|
||||
def test_service_matching_is_exact_for_overlapping_profile_names():
|
||||
assert update_cmd._gateway_service_matches_profile(
|
||||
"foo", "hermes-gateway-foo.service"
|
||||
@@ -186,7 +253,12 @@ def test_recovery_child_restarts_each_profile_with_a_fresh_main(monkeypatch):
|
||||
monkeypatch.setenv("_HERMES_GATEWAY", "1")
|
||||
result = recovery.restart_profiles(["default", "coder"], run=fake_run)
|
||||
|
||||
assert result == {"succeeded": ["coder", "default"], "failed": []}
|
||||
# No supervisor observations were possible → conservative labels only.
|
||||
assert result == {
|
||||
"verified": [],
|
||||
"relaunch_attempted": ["coder", "default"],
|
||||
"failed": [],
|
||||
}
|
||||
assert [call[0] for call in calls] == [
|
||||
[sys.executable, "-m", "hermes_cli.main", "-p", "coder", "gateway", "restart"],
|
||||
[sys.executable, "-m", "hermes_cli.main", "-p", "default", "gateway", "restart"],
|
||||
@@ -200,6 +272,52 @@ def test_recovery_child_restarts_each_profile_with_a_fresh_main(monkeypatch):
|
||||
assert "_HERMES_GATEWAY" not in kwargs["env"]
|
||||
|
||||
|
||||
def test_recovery_child_verifies_systemd_profiles_via_is_active(monkeypatch):
|
||||
recovery = importlib.import_module("hermes_cli.update_restart_recovery")
|
||||
monkeypatch.setattr(recovery.shutil, "which", lambda name: f"/bin/{name}")
|
||||
calls = []
|
||||
|
||||
def fake_run(argv, **kwargs):
|
||||
calls.append(argv)
|
||||
if argv[0].endswith("systemctl"):
|
||||
unit = argv[-1]
|
||||
active = unit == "hermes-gateway.service"
|
||||
return _Completed(0 if active else 3, stdout="active" if active else "inactive")
|
||||
return _Completed(0)
|
||||
|
||||
result = recovery.restart_profiles(
|
||||
["default", "coder"],
|
||||
supervisors={"default": "systemd", "coder": "launchd"},
|
||||
run=fake_run,
|
||||
)
|
||||
|
||||
assert result == {
|
||||
"verified": ["default"],
|
||||
"relaunch_attempted": ["coder"],
|
||||
"failed": [],
|
||||
}
|
||||
# The launchd profile must never be probed with systemctl.
|
||||
systemctl_units = [argv[-1] for argv in calls if argv[0].endswith("systemctl")]
|
||||
assert all("coder" not in unit for unit in systemctl_units)
|
||||
|
||||
|
||||
def test_recovery_child_treats_missing_systemctl_as_unverified(monkeypatch):
|
||||
recovery = importlib.import_module("hermes_cli.update_restart_recovery")
|
||||
monkeypatch.setattr(recovery.shutil, "which", lambda name: None)
|
||||
|
||||
result = recovery.restart_profiles(
|
||||
["default"],
|
||||
supervisors={"default": "systemd"},
|
||||
run=lambda *args, **kwargs: _Completed(0),
|
||||
)
|
||||
|
||||
assert result == {
|
||||
"verified": [],
|
||||
"relaunch_attempted": ["default"],
|
||||
"failed": [],
|
||||
}
|
||||
|
||||
|
||||
def test_recovery_child_reports_failed_profile_without_losing_successes():
|
||||
recovery = importlib.import_module("hermes_cli.update_restart_recovery")
|
||||
outcomes = iter((_Completed(1), _Completed(0)))
|
||||
@@ -208,7 +326,11 @@ def test_recovery_child_reports_failed_profile_without_losing_successes():
|
||||
["coder", "default"], run=lambda *args, **kwargs: next(outcomes)
|
||||
)
|
||||
|
||||
assert result == {"succeeded": ["default"], "failed": ["coder"]}
|
||||
assert result == {
|
||||
"verified": [],
|
||||
"relaunch_attempted": ["default"],
|
||||
"failed": ["coder"],
|
||||
}
|
||||
|
||||
|
||||
def test_recovery_payload_rejects_path_like_profile_ids():
|
||||
@@ -222,6 +344,23 @@ def test_recovery_payload_rejects_path_like_profile_ids():
|
||||
raise AssertionError("path-like profile id must be rejected")
|
||||
|
||||
|
||||
def test_recovery_payload_rejects_malformed_supervisors_map():
|
||||
recovery = importlib.import_module("hermes_cli.update_restart_recovery")
|
||||
|
||||
try:
|
||||
recovery._parse_payload(
|
||||
io.StringIO(
|
||||
json.dumps(
|
||||
{"profiles": ["default"], "supervisors": {"default": "sys/temd"}}
|
||||
)
|
||||
)
|
||||
)
|
||||
except ValueError as exc:
|
||||
assert "supervisors" in str(exc)
|
||||
else:
|
||||
raise AssertionError("malformed supervisors map must be rejected")
|
||||
|
||||
|
||||
def test_recovery_module_empty_payload_is_a_real_clean_process():
|
||||
result = subprocess.run(
|
||||
[sys.executable, "-m", "hermes_cli.update_restart_recovery", "--stdin"],
|
||||
@@ -232,4 +371,92 @@ def test_recovery_module_empty_payload_is_a_real_clean_process():
|
||||
)
|
||||
|
||||
assert result.returncode == 0
|
||||
assert json.loads(result.stdout) == {"succeeded": [], "failed": []}
|
||||
assert json.loads(result.stdout) == {
|
||||
"failed": [],
|
||||
"relaunch_attempted": [],
|
||||
"verified": [],
|
||||
}
|
||||
|
||||
|
||||
def test_recovery_module_end_to_end_in_a_real_fresh_process(tmp_path):
|
||||
"""E2E: the whole recovery protocol through a genuinely fresh interpreter.
|
||||
|
||||
A ``sitecustomize`` shim in the child's ``PYTHONPATH`` intercepts the
|
||||
grandchild ``hermes_cli.main … gateway restart`` invocations (recording
|
||||
them and returning rc 0) and answers ``systemctl --user is-active`` with
|
||||
``active`` only for the default profile's unit. Everything else — stdin
|
||||
payload parsing, profile ordering, environment scrubbing, verification
|
||||
classification, JSON output, and exit code — runs the real module code in
|
||||
a real new process, exactly as the aborted updater would spawn it.
|
||||
"""
|
||||
ledger = tmp_path / "grandchild_calls.jsonl"
|
||||
shim = textwrap.dedent(
|
||||
f"""
|
||||
import json
|
||||
import shutil
|
||||
import subprocess
|
||||
|
||||
_real_run = subprocess.run
|
||||
_real_which = shutil.which
|
||||
_LEDGER = {str(ledger)!r}
|
||||
|
||||
|
||||
def _shim_which(name, *args, **kwargs):
|
||||
if name == "systemctl":
|
||||
return "/usr/bin/systemctl"
|
||||
return _real_which(name, *args, **kwargs)
|
||||
|
||||
|
||||
shutil.which = _shim_which
|
||||
|
||||
|
||||
def _shim_run(argv, *args, **kwargs):
|
||||
argv_list = list(argv)
|
||||
if "hermes_cli.main" in argv_list:
|
||||
with open(_LEDGER, "a", encoding="utf-8") as fh:
|
||||
fh.write(json.dumps(argv_list) + "\\n")
|
||||
return subprocess.CompletedProcess(argv_list, 0, "", "")
|
||||
if argv_list and str(argv_list[0]).endswith("systemctl"):
|
||||
unit = argv_list[-1]
|
||||
if unit == "hermes-gateway.service":
|
||||
return subprocess.CompletedProcess(argv_list, 0, "active\\n", "")
|
||||
return subprocess.CompletedProcess(argv_list, 3, "inactive\\n", "")
|
||||
return _real_run(argv, *args, **kwargs)
|
||||
|
||||
|
||||
subprocess.run = _shim_run
|
||||
"""
|
||||
)
|
||||
(tmp_path / "sitecustomize.py").write_text(shim, encoding="utf-8")
|
||||
|
||||
import os
|
||||
|
||||
env = os.environ.copy()
|
||||
env["PYTHONPATH"] = str(tmp_path) + os.pathsep + env.get("PYTHONPATH", "")
|
||||
env["_HERMES_GATEWAY"] = "1" # must be scrubbed before the grandchild runs
|
||||
|
||||
result = subprocess.run(
|
||||
[sys.executable, "-m", "hermes_cli.update_restart_recovery", "--stdin"],
|
||||
input=json.dumps(
|
||||
{
|
||||
"profiles": ["default", "coder"],
|
||||
"supervisors": {"default": "systemd", "coder": "launchd"},
|
||||
}
|
||||
),
|
||||
capture_output=True,
|
||||
text=True,
|
||||
check=False,
|
||||
env=env,
|
||||
timeout=120,
|
||||
)
|
||||
|
||||
assert result.returncode == 0, result.stderr
|
||||
assert json.loads(result.stdout) == {
|
||||
"failed": [],
|
||||
"relaunch_attempted": ["coder"],
|
||||
"verified": ["default"],
|
||||
}
|
||||
restarts = [json.loads(line) for line in ledger.read_text().splitlines()]
|
||||
assert [argv[argv.index("-p") + 1] for argv in restarts] == ["coder", "default"]
|
||||
for argv in restarts:
|
||||
assert argv[-2:] == ["gateway", "restart"]
|
||||
|
||||
@@ -1751,7 +1751,7 @@ Pulls the latest `hermes-agent` code and reinstalls dependencies in the managed
|
||||
Additional behavior:
|
||||
|
||||
- **Gateway restart.** After a successful update, Hermes attempts to restart all running gateway profiles automatically so they pick up the new code. Use `hermes gateway restart` when you want to restart a gateway without applying an update.
|
||||
- **Restart-phase recovery.** If the in-process restart phase aborts while importing the freshly pulled tree, supervised gateway profiles are retried through a clean Python process. Manual gateways are not killed without a relaunch authority; they remain in the incomplete-update report with the exact restart command.
|
||||
- **Restart-phase recovery.** If the in-process restart phase aborts while importing the freshly pulled tree, supervised gateway profiles are retried through a clean Python process. Only restarts independently confirmed by systemd (`systemctl --user is-active`) are reported as verified; a relaunch that merely exited 0 is recorded as `relaunch_attempted` and still fails the update conservatively. Manual gateways and serve/dashboard runtimes are never killed without a relaunch authority; they are recorded as skipped with a reason and remain in the incomplete-update report with the exact restart command.
|
||||
- **Update receipts + fleet version check.** Every run writes a machine-readable receipt to `~/.hermes/logs/update_receipts/` (pre-update fleet plan, steps, skips with reasons, restart outcome; `latest.json` points at the newest). After the restart phase the updater verifies each live gateway's running code against the updated checkout and prints a per-profile version matrix; a gateway still on pre-update code fails the update (exit 1) with the exact restart command.
|
||||
- **Local source changes.** For git installs, dirty tracked files and untracked files are auto-stashed before branch checkout or pull (`git stash push --include-untracked`). Interactive terminal updates ask before restoring the stash. Non-interactive updates restore it by default; set `updates.non_interactive_local_changes: discard` only on managed installs where local source edits should be thrown away after a successful pull. If stash restore conflicts or the pull fails, the stash is left in place for manual recovery.
|
||||
- **npm lockfile churn.** Before stashing or switching branches, Hermes makes a best-effort cleanup of tracked `package-lock.json` diffs produced by npm install/build steps. Commit or manually stash intentional lockfile edits before running `hermes update`.
|
||||
|
||||
Reference in New Issue
Block a user