From 46eaefee4ebf549011c31f0245319ea951da720c Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 21:03:09 -0700 Subject: [PATCH] =?UTF-8?q?refactor(hermes=5Fcli):=20keepalive/process=5Fi?= =?UTF-8?q?dentity=20=E2=80=94=20compact=20signatures,=20fold=20guard=20ch?= =?UTF-8?q?ains=20(699->660=20LOC)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- hermes_cli/nous_auth_keepalive.py | 59 ++++++++----------------------- hermes_cli/process_identity.py | 18 +++------- 2 files changed, 19 insertions(+), 58 deletions(-) diff --git a/hermes_cli/nous_auth_keepalive.py b/hermes_cli/nous_auth_keepalive.py index cbbc5121aa..ae5538dc38 100644 --- a/hermes_cli/nous_auth_keepalive.py +++ b/hermes_cli/nous_auth_keepalive.py @@ -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") diff --git a/hermes_cli/process_identity.py b/hermes_cli/process_identity.py index 95b975e918..2c666c7a08 100644 --- a/hermes_cli/process_identity.py +++ b/hermes_cli/process_identity.py @@ -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