fix(gateway): read reconnect_attention_after from the bound profile's config, not the env bridge
Follow-up to treatux's call-time read: resolving `_float_env(...)` per call still reads `os.environ`, which under a multiplexed gateway holds the LAUNCH profile's bridged value — profile B in `_profile_runtime_scope` kept escalating at A's threshold. The env bridge entry goes away (no other consumer); `_reconnect_attention_after_secs()` reads `agent.reconnect_attention_after` via `load_config_readonly()` so the value follows the scope bound at call time and a config edit needs no restart (mtime-keyed cache). The secondary-profile reconnect loop (`_run_secondary_profile_reconnect`) never escalated at all: it now carries a queue entry and flags `<profile>:<platform>` NEEDS_ATTENTION inside its own profile scope, mirroring the primary watcher. Closes #115635 Co-authored-by: treatux <290971332+treatux@users.noreply.github.com>
This commit is contained in:
@@ -1900,8 +1900,6 @@ _AGENT_ENV_BRIDGE = {
|
||||
"gateway_timeout_warning": "HERMES_AGENT_TIMEOUT_WARNING",
|
||||
"gateway_notify_interval": "HERMES_AGENT_NOTIFY_INTERVAL",
|
||||
"session_stall_timeout": "HERMES_SESSION_STALL_TIMEOUT",
|
||||
# Internal bridge only — config.yaml (agent.reconnect_attention_after) is the documented setting.
|
||||
"reconnect_attention_after": "HERMES_RECONNECT_ATTENTION_AFTER_SECONDS",
|
||||
"restart_drain_timeout": "HERMES_RESTART_DRAIN_TIMEOUT",
|
||||
"cron_drain_timeout": "HERMES_CRON_DRAIN_TIMEOUT",
|
||||
"gateway_auto_continue_freshness": "HERMES_AUTO_CONTINUE_FRESHNESS",
|
||||
@@ -3223,19 +3221,30 @@ async def _dispose_unused_adapter(adapter: "BasePlatformAdapter | None") -> None
|
||||
# Max seconds between platform reconnect retries (primary watcher and secondary profiles share it).
|
||||
_RECONNECT_BACKOFF_CAP = 300
|
||||
|
||||
# Seconds continuously in the reconnect queue before NEEDS_ATTENTION. Retrying never stops (transient
|
||||
# outages must self-heal); this only makes a permanently-failing loop loud. 0 disables.
|
||||
|
||||
|
||||
def _reconnect_backoff(attempt: int) -> int:
|
||||
"""Exponential reconnect backoff: 30s, 60s, 120s, ... capped at 5 min."""
|
||||
return min(30 * (2 ** (attempt - 1)), _RECONNECT_BACKOFF_CAP)
|
||||
|
||||
|
||||
def _reconnect_attention_after_secs() -> float:
|
||||
"""``agent.reconnect_attention_after`` of the profile whose scope is bound at call time (the launch
|
||||
profile's when unbound). Seconds continuously in the reconnect queue before NEEDS_ATTENTION; retrying
|
||||
never stops (transient outages must self-heal), this only makes a permanently-failing loop loud.
|
||||
Non-positive disables. Read per call, never cached: one process serves many profiles and a config
|
||||
edit must not need a gateway restart (#115635)."""
|
||||
from hermes_cli.config import load_config_readonly
|
||||
agent_cfg = load_config_readonly().get("agent")
|
||||
raw = agent_cfg.get("reconnect_attention_after") if isinstance(agent_cfg, dict) else None
|
||||
try:
|
||||
return float(raw)
|
||||
except (TypeError, ValueError):
|
||||
return float(_DEFAULT_CONFIG["agent"]["reconnect_attention_after"])
|
||||
|
||||
|
||||
def _reconnect_needs_attention(info: dict, now: float) -> bool:
|
||||
"""True when a reconnect-queue entry has waited long enough for NEEDS_ATTENTION.
|
||||
``queued_at`` is re-stamped on each (re)entry, so only *continuous* failure escalates."""
|
||||
threshold = _float_env("HERMES_RECONNECT_ATTENTION_AFTER_SECONDS", 7200)
|
||||
threshold = _reconnect_attention_after_secs()
|
||||
if threshold <= 0:
|
||||
return False # escalation disabled
|
||||
queued_at = info.get("queued_at")
|
||||
|
||||
@@ -657,8 +657,12 @@ class GatewayAdapterLifecycleMixin:
|
||||
if not await _idle(10): # re-check every 10 seconds
|
||||
return
|
||||
|
||||
def _flag_reconnect_needs_attention(self, platform, info: dict, now: float) -> None:
|
||||
"""Flag NEEDS_ATTENTION (once) past the threshold — a signal, NOT a circuit breaker."""
|
||||
def _flag_reconnect_needs_attention(
|
||||
self, platform, info: dict, now: float, *, status_key: Optional[str] = None
|
||||
) -> None:
|
||||
"""Flag NEEDS_ATTENTION (once) past the threshold — a signal, NOT a circuit breaker. The threshold
|
||||
is the bound profile's ``agent.reconnect_attention_after``: secondaries call this inside their
|
||||
``_profile_runtime_scope`` with their ``<profile>:<platform>`` status key."""
|
||||
from gateway.run import _reconnect_needs_attention
|
||||
if info.get("attention_flagged") or not _reconnect_needs_attention(info, now):
|
||||
return
|
||||
@@ -668,10 +672,10 @@ class GatewayAdapterLifecycleMixin:
|
||||
"%s has been failing/reconnecting continuously for %.1f hours (%d attempts) — flagging "
|
||||
"NEEDS_ATTENTION. Retries continue, but this usually means a permanent problem (revoked "
|
||||
"credentials, missing intents, broken sidecar). Check `hermes status` / `/platform list`.",
|
||||
platform.value, queued_for / 3600.0, info.get("attempts", 0),
|
||||
status_key or platform.value, queued_for / 3600.0, info.get("attempts", 0),
|
||||
)
|
||||
self._update_platform_runtime_status(
|
||||
platform.value, platform_state="retrying", needs_attention=True,
|
||||
status_key or platform.value, platform_state="retrying", needs_attention=True,
|
||||
retrying_since=(datetime.now(timezone.utc) - timedelta(seconds=queued_for)).isoformat(),
|
||||
)
|
||||
|
||||
@@ -1191,8 +1195,10 @@ class GatewayAdapterLifecycleMixin:
|
||||
|
||||
async def _run_secondary_profile_reconnect(self, profile_name: str, platform: Platform) -> None:
|
||||
"""Reconnect a retryable secondary adapter under its own profile scope."""
|
||||
from gateway.run import _reconnect_backoff
|
||||
from gateway.run import _profile_runtime_scope, _reconnect_backoff
|
||||
attempts = 0
|
||||
# Same escalation shape as the primary queue entry; ``queued_at`` is this task's start.
|
||||
queue_info = {"queued_at": time.monotonic(), "attempts": 0}
|
||||
current_task = asyncio.current_task()
|
||||
try:
|
||||
while self._running:
|
||||
@@ -1231,6 +1237,13 @@ class GatewayAdapterLifecycleMixin:
|
||||
if not self._running:
|
||||
return
|
||||
attempts += 1
|
||||
queue_info["attempts"] = attempts
|
||||
profile_home = self._profile_home_or_none(profile_name)
|
||||
# The attempt above already hydrated this profile's secret sources off-loop.
|
||||
with self._scope_or_null(
|
||||
functools.partial(_profile_runtime_scope, hydrate_secrets=False), profile_home):
|
||||
self._flag_reconnect_needs_attention(
|
||||
platform, queue_info, time.monotonic(), status_key=f"{profile_name}:{platform.value}")
|
||||
backoff = _reconnect_backoff(attempts)
|
||||
logger.info(
|
||||
"Secondary %s reconnect retry in %ds (profile: %s)", platform.value, backoff, profile_name
|
||||
|
||||
@@ -1,107 +1,83 @@
|
||||
"""Profile-scope leak regression for reconnect attention threshold.
|
||||
"""``agent.reconnect_attention_after`` is read from the profile whose scope is bound at call time.
|
||||
|
||||
Issue: _RECONNECT_ATTENTION_AFTER_SECONDS was a module-level constant that
|
||||
captured os.environ at import time. Under multiplex, switching profiles and
|
||||
re-bridging config had no effect — the threshold stayed locked to the launch
|
||||
profile's value.
|
||||
|
||||
Test: two HERMES_HOME dirs (A -> B) prove the threshold moves with the
|
||||
active profile after the fix.
|
||||
Regression for #115635: the threshold was a module constant frozen from ``os.environ`` at import,
|
||||
so a multiplexed secondary escalated at the LAUNCH profile's value and a live config edit needed a
|
||||
gateway restart. Two homes, A -> B -> A, per AGENTS.md § "One process may serve many profiles".
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import subprocess
|
||||
import sys
|
||||
import textwrap
|
||||
from pathlib import Path
|
||||
import asyncio
|
||||
import time
|
||||
|
||||
import pytest
|
||||
import yaml
|
||||
|
||||
PROJECT_ROOT = Path(__file__).resolve().parents[2]
|
||||
import gateway.run as gateway_run
|
||||
from gateway.config import Platform
|
||||
from gateway.run import GatewayRunner
|
||||
|
||||
|
||||
def _write_config(home: Path, agent_cfg: dict | None = None) -> None:
|
||||
cfg: dict = {}
|
||||
if agent_cfg:
|
||||
cfg["agent"] = agent_cfg
|
||||
(home / "config.yaml").write_text(yaml.safe_dump(cfg), encoding="utf-8")
|
||||
def _home_with_threshold(root, name, seconds):
|
||||
home = root / name
|
||||
home.mkdir()
|
||||
(home / "config.yaml").write_text(
|
||||
yaml.safe_dump({"agent": {"reconnect_attention_after": seconds}}), encoding="utf-8")
|
||||
return home
|
||||
|
||||
|
||||
def test_reconnect_attention_respects_profile_switch(tmp_path: Path) -> None:
|
||||
"""Bridging a different profile's config must update the reconnect threshold.
|
||||
def test_threshold_follows_bound_profile_scope_a_b_a(tmp_path, monkeypatch):
|
||||
home_a = _home_with_threshold(tmp_path, "a", 10)
|
||||
home_b = _home_with_threshold(tmp_path, "b", 100000)
|
||||
monkeypatch.setenv("HERMES_HOME", str(home_a))
|
||||
now = time.monotonic()
|
||||
queued_20s_ago = {"queued_at": now - 20}
|
||||
|
||||
Before the fix, _RECONNECT_ATTENTION_AFTER_SECONDS was cached at import time,
|
||||
so profile B inherited profile A's threshold silently.
|
||||
"""
|
||||
home_a = tmp_path / "profile_a"
|
||||
home_b = tmp_path / "profile_b"
|
||||
home_a.mkdir()
|
||||
home_b.mkdir()
|
||||
assert gateway_run._reconnect_needs_attention(dict(queued_20s_ago), now) is True
|
||||
with gateway_run._profile_runtime_scope(home_b, hydrate_secrets=False):
|
||||
assert gateway_run._reconnect_needs_attention(dict(queued_20s_ago), now) is False
|
||||
assert gateway_run._reconnect_needs_attention(dict(queued_20s_ago), now) is True
|
||||
|
||||
_write_config(home_a, agent_cfg={"reconnect_attention_after": 10})
|
||||
_write_config(home_b, agent_cfg={"reconnect_attention_after": 20})
|
||||
# Live edit of the bound profile's config takes effect on the next call, no restart.
|
||||
(home_a / "config.yaml").write_text(
|
||||
yaml.safe_dump({"agent": {"reconnect_attention_after": 100000}}), encoding="utf-8")
|
||||
assert gateway_run._reconnect_needs_attention(dict(queued_20s_ago), now) is False
|
||||
|
||||
script = textwrap.dedent(
|
||||
f"""
|
||||
import os, sys, time
|
||||
sys.path.insert(0, {str(PROJECT_ROOT)!r})
|
||||
|
||||
from pathlib import Path
|
||||
from gateway import run
|
||||
from gateway.run import _bridge_config_to_env, _load_bridge_config
|
||||
@pytest.mark.asyncio
|
||||
async def test_secondary_reconnect_loop_escalates_under_own_profile(tmp_path, monkeypatch):
|
||||
"""A secondary profile's reconnect loop flags ``<profile>:<platform>`` NEEDS_ATTENTION at ITS
|
||||
threshold — the launch profile's (100000 here) must not suppress it."""
|
||||
launch = _home_with_threshold(tmp_path, "launch", 100000)
|
||||
secondary = _home_with_threshold(tmp_path, "sec", 0.01)
|
||||
monkeypatch.setenv("HERMES_HOME", str(launch))
|
||||
|
||||
# Simulate serving profile A first
|
||||
cfg_a = _load_bridge_config(Path({str(home_a / 'config.yaml')!r}))
|
||||
_bridge_config_to_env(cfg_a)
|
||||
runner = object.__new__(GatewayRunner)
|
||||
runner._running = True
|
||||
runner._profile_adapters = {}
|
||||
runner._profile_failed_platforms = {}
|
||||
writes = []
|
||||
monkeypatch.setattr(runner, "_update_platform_runtime_status", lambda key, **kw: writes.append((key, kw)))
|
||||
|
||||
# 15 s > A's threshold of 10 -> should flag
|
||||
info_a = {{"queued_at": time.monotonic() - 15}}
|
||||
assert run._reconnect_needs_attention(info_a, time.monotonic()) is True, \
|
||||
"A threshold 10: 15s should flag"
|
||||
class _RetryableAdapter:
|
||||
has_fatal_error = True
|
||||
fatal_error_retryable = True
|
||||
|
||||
# 5 s < A's threshold of 10 -> should not flag
|
||||
info_a_ok = {{"queued_at": time.monotonic() - 5}}
|
||||
assert run._reconnect_needs_attention(info_a_ok, time.monotonic()) is False, \
|
||||
"A threshold 10: 5s should not flag"
|
||||
async def failing_attempt(profile_name, platform):
|
||||
return _RetryableAdapter(), False
|
||||
|
||||
# Switch to profile B (same process, multiplex scenario)
|
||||
cfg_b = _load_bridge_config(Path({str(home_b / 'config.yaml')!r}))
|
||||
_bridge_config_to_env(cfg_b)
|
||||
async def noop_disconnect(adapter, platform):
|
||||
return None
|
||||
|
||||
# 15 s < B's threshold of 20 -> must NOT flag after fix
|
||||
info_b = {{"queued_at": time.monotonic() - 15}}
|
||||
result = run._reconnect_needs_attention(info_b, time.monotonic())
|
||||
print(f"RESULT={{result}}")
|
||||
"""
|
||||
)
|
||||
monkeypatch.setattr(runner, "_secondary_reconnect_attempt", failing_attempt)
|
||||
monkeypatch.setattr(runner, "_safe_adapter_disconnect", noop_disconnect)
|
||||
monkeypatch.setattr(runner, "_profile_home_or_none", lambda name: secondary)
|
||||
monkeypatch.setattr(gateway_run, "_reconnect_backoff", lambda attempts: 0.02)
|
||||
|
||||
env = dict(os.environ)
|
||||
env["HERMES_HOME"] = str(home_a)
|
||||
# Keep interpreter paths required by stdlib / platform detection
|
||||
for k in (
|
||||
"PATH", "PYTHONPATH", "VIRTUAL_ENV", "HOME", "USERPROFILE",
|
||||
"HOMEDRIVE", "HOMEPATH", "LOCALAPPDATA", "APPDATA",
|
||||
"SYSTEMROOT", "TEMP", "TMP",
|
||||
):
|
||||
if k in os.environ and k not in env:
|
||||
env[k] = os.environ[k]
|
||||
task = asyncio.create_task(runner._run_secondary_profile_reconnect("sec", Platform.DISCORD))
|
||||
deadline = time.monotonic() + 5
|
||||
while time.monotonic() < deadline and not any(kw.get("needs_attention") for _k, kw in writes):
|
||||
await asyncio.sleep(0.02)
|
||||
runner._running = False
|
||||
await asyncio.wait_for(task, timeout=2)
|
||||
|
||||
proc = subprocess.run(
|
||||
[sys.executable, "-c", script],
|
||||
env=env,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=60,
|
||||
)
|
||||
if proc.returncode != 0:
|
||||
pytest.fail(
|
||||
f"Subscript failed (rc={proc.returncode})\n"
|
||||
f"stderr:\n{proc.stderr}\nstdout:\n{proc.stdout}"
|
||||
)
|
||||
assert "RESULT=False" in proc.stdout, (
|
||||
f"Expected False after B bridge (threshold 20, elapsed 15), got:\n"
|
||||
f"{proc.stdout}"
|
||||
)
|
||||
flagged = [(key, kw) for key, kw in writes if kw.get("needs_attention")]
|
||||
assert [key for key, _kw in flagged] == ["sec:discord"], writes
|
||||
assert flagged[0][1]["platform_state"] == "retrying" and flagged[0][1].get("retrying_since")
|
||||
|
||||
Reference in New Issue
Block a user