refactor(hermes_cli): keepalive/process_identity — compact signatures, fold guard chains (699->660 LOC)
This commit is contained in:
@@ -94,7 +94,7 @@ def _observed_lifetime_seconds() -> Optional[int]:
|
||||
continue
|
||||
if value > 0:
|
||||
lifetimes.append(value)
|
||||
return min(lifetimes) if lifetimes else None
|
||||
return min(lifetimes, default=None)
|
||||
|
||||
|
||||
def _tick_seconds(configured_interval: int, lifetime: Optional[int]) -> int:
|
||||
@@ -113,18 +113,10 @@ def _refresh_horizon_seconds(tick_seconds: int, floor_seconds: int) -> int:
|
||||
|
||||
|
||||
def _entry_state(entry: object) -> dict:
|
||||
return {
|
||||
"agent_key": getattr(entry, "agent_key", None),
|
||||
"agent_key_expires_at": getattr(entry, "agent_key_expires_at", None),
|
||||
"scope": getattr(entry, "scope", None),
|
||||
}
|
||||
return {k: getattr(entry, k, None) for k in ("agent_key", "agent_key_expires_at", "scope")}
|
||||
|
||||
|
||||
def _refresh_selected_pool_entry(
|
||||
*,
|
||||
min_key_ttl_seconds: int,
|
||||
min_access_ttl_seconds: Optional[int] = None,
|
||||
) -> Optional[bool]:
|
||||
def _refresh_selected_pool_entry(*, min_key_ttl_seconds: int, min_access_ttl_seconds: Optional[int] = None) -> Optional[bool]:
|
||||
"""Refresh the current pool entry when stale. True = usable/refreshed; False = pool exists but
|
||||
no usable entry; None = no Nous pool.
|
||||
"""
|
||||
@@ -135,19 +127,15 @@ def _refresh_selected_pool_entry(
|
||||
except Exception as exc:
|
||||
logger.debug("Nous auth keepalive: credential pool unavailable: %s", exc)
|
||||
return None
|
||||
|
||||
if not pool or not pool.has_credentials():
|
||||
return None
|
||||
|
||||
try:
|
||||
entry = pool.select()
|
||||
except Exception as exc:
|
||||
logger.debug("Nous auth keepalive: credential pool selection failed: %s", exc)
|
||||
return False
|
||||
|
||||
if entry is None:
|
||||
return False
|
||||
|
||||
if min_access_ttl_seconds is None:
|
||||
min_access_ttl_seconds = ACCESS_TOKEN_REFRESH_SKEW_SECONDS
|
||||
access_expiring = _is_expiring(getattr(entry, "expires_at", None), min_access_ttl_seconds)
|
||||
@@ -160,24 +148,17 @@ def _refresh_selected_pool_entry(
|
||||
|
||||
|
||||
def refresh_nous_auth_keepalive_once(
|
||||
*,
|
||||
min_key_ttl_seconds: int = NOUS_INVOKE_JWT_MIN_TTL_SECONDS,
|
||||
min_access_ttl_seconds: Optional[int] = None,
|
||||
timeout_seconds: Optional[float] = None,
|
||||
*, min_key_ttl_seconds: int = NOUS_INVOKE_JWT_MIN_TTL_SECONDS,
|
||||
min_access_ttl_seconds: Optional[int] = None, timeout_seconds: Optional[float] = None,
|
||||
) -> bool:
|
||||
"""Refresh Nous auth once if credentials are configured."""
|
||||
min_key_ttl_seconds = max(60, int(min_key_ttl_seconds))
|
||||
|
||||
"""Refresh Nous auth once if credentials are configured (pool entry first, then singleton state)."""
|
||||
pool_result = _refresh_selected_pool_entry(
|
||||
min_key_ttl_seconds=min_key_ttl_seconds, min_access_ttl_seconds=min_access_ttl_seconds
|
||||
min_key_ttl_seconds=max(60, int(min_key_ttl_seconds)), min_access_ttl_seconds=min_access_ttl_seconds
|
||||
)
|
||||
if pool_result is not None:
|
||||
return pool_result
|
||||
|
||||
state = get_provider_auth_state("nous")
|
||||
if not state:
|
||||
if not get_provider_auth_state("nous"):
|
||||
return False
|
||||
|
||||
try:
|
||||
resolve_nous_runtime_credentials(timeout_seconds=_timeout_seconds(timeout_seconds))
|
||||
logger.debug("Nous auth keepalive: refreshed singleton auth state")
|
||||
@@ -191,16 +172,11 @@ def refresh_nous_auth_keepalive_once(
|
||||
|
||||
|
||||
def _keepalive_loop(
|
||||
stop_event: threading.Event,
|
||||
*,
|
||||
interval_seconds: int,
|
||||
initial_delay_seconds: int,
|
||||
min_key_ttl_seconds: int,
|
||||
timeout_seconds: Optional[float],
|
||||
stop_event: threading.Event, *, interval_seconds: int, initial_delay_seconds: int,
|
||||
min_key_ttl_seconds: int, timeout_seconds: Optional[float],
|
||||
) -> None:
|
||||
if initial_delay_seconds > 0 and stop_event.wait(initial_delay_seconds):
|
||||
return
|
||||
|
||||
while not stop_event.is_set():
|
||||
# Re-read each pass: the lifetime changes with account/plan/policy; caching it would go
|
||||
# stale in exactly the case the keepalive exists to cover.
|
||||
@@ -213,34 +189,27 @@ def _keepalive_loop(
|
||||
|
||||
|
||||
def start_nous_auth_keepalive(
|
||||
*,
|
||||
interval_seconds: Optional[int] = None,
|
||||
*, interval_seconds: Optional[int] = None,
|
||||
initial_delay_seconds: int = NOUS_AUTH_KEEPALIVE_INITIAL_DELAY_SECONDS,
|
||||
min_key_ttl_seconds: int = NOUS_INVOKE_JWT_MIN_TTL_SECONDS,
|
||||
timeout_seconds: Optional[float] = None,
|
||||
min_key_ttl_seconds: int = NOUS_INVOKE_JWT_MIN_TTL_SECONDS, timeout_seconds: Optional[float] = None,
|
||||
) -> Optional[threading.Thread]:
|
||||
"""Start the process-wide Nous auth keepalive thread."""
|
||||
"""Start the process-wide Nous auth keepalive thread (idempotent; None when disabled)."""
|
||||
interval_seconds = _interval_seconds(interval_seconds)
|
||||
if interval_seconds <= 0:
|
||||
return None
|
||||
|
||||
global _keepalive_thread
|
||||
with _keepalive_lock:
|
||||
if _keepalive_thread is not None and _keepalive_thread.is_alive():
|
||||
return _keepalive_thread
|
||||
|
||||
_keepalive_stop.clear()
|
||||
_keepalive_thread = threading.Thread(
|
||||
target=_keepalive_loop,
|
||||
args=(_keepalive_stop,),
|
||||
target=_keepalive_loop, args=(_keepalive_stop,), daemon=True, name="nous-auth-keepalive",
|
||||
kwargs={
|
||||
"interval_seconds": int(interval_seconds),
|
||||
"initial_delay_seconds": max(0, int(initial_delay_seconds)),
|
||||
"min_key_ttl_seconds": max(60, int(min_key_ttl_seconds)),
|
||||
"timeout_seconds": timeout_seconds,
|
||||
},
|
||||
daemon=True,
|
||||
name="nous-auth-keepalive",
|
||||
)
|
||||
_keepalive_thread.start()
|
||||
logger.debug("Nous auth keepalive started")
|
||||
|
||||
@@ -159,9 +159,7 @@ def _read_ledger(path: Path) -> Optional[list[dict]]:
|
||||
parsed = json.loads(text)
|
||||
except (ValueError, TypeError):
|
||||
return None
|
||||
if not isinstance(parsed, list):
|
||||
return None
|
||||
return [e for e in parsed if isinstance(e, dict)]
|
||||
return [e for e in parsed if isinstance(e, dict)] if isinstance(parsed, list) else None
|
||||
|
||||
|
||||
def _read_ledger_or_quarantine(path: Path) -> Optional[list[dict]]:
|
||||
@@ -290,9 +288,7 @@ def register_child(pid: int, purpose: str, *, project_root: Optional[Path] = Non
|
||||
pid = int(pid)
|
||||
except (TypeError, ValueError):
|
||||
return False
|
||||
if pid <= 0:
|
||||
return False
|
||||
child_create = _process_create_time(pid)
|
||||
child_create = _process_create_time(pid) if pid > 0 else None
|
||||
if child_create is None:
|
||||
return False
|
||||
entry = _new_entry(pid, child_create, purpose, project_root, os.getpid(), _process_create_time())
|
||||
@@ -321,7 +317,7 @@ def ledger_entries(*, project_root: Optional[Path] = None) -> list[dict]:
|
||||
|
||||
|
||||
def spawner_is_dead(entry: dict) -> Optional[bool]:
|
||||
"""Is the recorded spawner of this entry provably gone?"""
|
||||
"""Is the recorded spawner of this entry provably gone? ``None`` when unrecorded/unprovable."""
|
||||
spawner_pid = entry.get("spawner_pid")
|
||||
if not isinstance(spawner_pid, int) or spawner_pid <= 0:
|
||||
return None
|
||||
@@ -344,10 +340,8 @@ def reap_orphaned_mcp_helpers(*, project_root: Optional[Path] = None, kill_fn=No
|
||||
own_pid = os.getpid()
|
||||
for entry in entries:
|
||||
try:
|
||||
if entry.get("purpose") != "mcp-helper":
|
||||
continue
|
||||
pid = entry.get("pid")
|
||||
if not isinstance(pid, int) or pid <= 0 or pid == own_pid:
|
||||
if entry.get("purpose") != "mcp-helper" or not isinstance(pid, int) or pid <= 0 or pid == own_pid:
|
||||
continue
|
||||
if spawner_is_dead(entry) is not True:
|
||||
continue # live or unprovable spawner → never touch
|
||||
@@ -425,9 +419,7 @@ def attach_self_to_kill_on_close_job() -> bool:
|
||||
info.BasicLimitInformation.LimitFlags = (
|
||||
JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE | JOB_OBJECT_LIMIT_BREAKAWAY_OK | JOB_OBJECT_LIMIT_SILENT_BREAKAWAY_OK
|
||||
)
|
||||
ok = kernel32.SetInformationJobObject(
|
||||
job, JobObjectExtendedLimitInformation, ctypes.byref(info), ctypes.sizeof(info)
|
||||
)
|
||||
ok = kernel32.SetInformationJobObject(job, JobObjectExtendedLimitInformation, ctypes.byref(info), ctypes.sizeof(info))
|
||||
if not ok or not kernel32.AssignProcessToJobObject(job, kernel32.GetCurrentProcess()):
|
||||
kernel32.CloseHandle(job)
|
||||
return False
|
||||
|
||||
Reference in New Issue
Block a user