fix(gateway): hot-serve reaches pooled Desktop backends; deleted profiles leave no stale runtime entries
- PUT /api/messaging/platforms on a pooled `hermes --profile X serve` arrives without ?profile= (Desktop local topology, #109088): resolve the hot-serve target from the process's own profile so the multiplexer is pinged and the UI skips the restart banner. - A profile deleted while the reconcile lock was held by its own adapter connect was recorded back into served_profiles; re-check the live set before recording. - Drop a deleted profile's `<name>:<platform>` runtime-status entries instead of leaving them as `stopped`.
This commit is contained in:
@@ -133,6 +133,13 @@ class GatewayProfileReconcileMixin:
|
||||
logger.info("[MULTIPLEX] Re-scanned profile '%s' after config/.env change (%s adapter(s) connected)", name, connected)
|
||||
result["rescanned"].append(name)
|
||||
self._served_profile_signatures = sigs
|
||||
# A profile deleted while an adapter above was still connecting must not be recorded back
|
||||
# (the deleter's signal timed out against this lock and rmtree already ran).
|
||||
live_now = {str(name) for name, _home in _multiplex_profile_homes(self.config)}
|
||||
for name in [n for n in current if n not in live_now and n != active]:
|
||||
await self._unserve_profile(name, current.pop(name))
|
||||
result["removed"].append(name)
|
||||
added = [n for n in added if n != name]
|
||||
self._record_served_profiles(active, list(current.items()))
|
||||
if added:
|
||||
await self._after_profiles_added([(n, current[n]) for n in added])
|
||||
@@ -180,7 +187,8 @@ class GatewayProfileReconcileMixin:
|
||||
adapters = (getattr(self, "_profile_adapters", None) or {}).pop(name, None) or {}
|
||||
for platform, adapter in list(adapters.items()):
|
||||
await self._bounded_adapter_teardown(adapter, platform, profile=name)
|
||||
_write_runtime_status_quiet(platform=f"{name}:{platform.value}", platform_state="stopped")
|
||||
# Its ``<name>:<platform>`` runtime entries describe a profile that no longer exists.
|
||||
_write_runtime_status_quiet(drop_profile_platforms=name)
|
||||
for attr in ("pairing_stores", "_busy_text_modes_by_profile", "_busy_input_modes_by_profile"):
|
||||
store = getattr(self, attr, None)
|
||||
if isinstance(store, dict):
|
||||
|
||||
@@ -805,19 +805,23 @@ def write_runtime_status(
|
||||
error_code: Any = _UNSET, error_message: Any = _UNSET, needs_attention: Any = _UNSET,
|
||||
retrying_since: Any = _UNSET, served_profiles: Any = _UNSET, session_store: Any = _UNSET,
|
||||
ingress_url: Any = _UNSET, clear_profile_platforms: bool = False,
|
||||
drop_profile_platforms: Optional[str] = None,
|
||||
) -> None:
|
||||
"""Persist gateway runtime health information for diagnostics/status."""
|
||||
"""Persist gateway runtime health information for diagnostics/status. ``drop_profile_platforms``
|
||||
removes one deleted profile's ``<profile>:<platform>`` entries (hot unroute)."""
|
||||
path = _get_runtime_status_path()
|
||||
payload = _read_json_file(path) or _build_runtime_status_record()
|
||||
previous_payload = copy.deepcopy(payload)
|
||||
current_record = _build_pid_record()
|
||||
payload.setdefault("platforms", {})
|
||||
if clear_profile_platforms:
|
||||
if clear_profile_platforms or drop_profile_platforms:
|
||||
# Secondary-profile entries are keyed ``<profile>:<platform>``. A fresh process must not
|
||||
# inherit them or /api/status stays degraded until every old adapter re-emits.
|
||||
platforms = payload["platforms"] if isinstance(payload["platforms"], dict) else {}
|
||||
drop_prefix = f"{drop_profile_platforms}:" if drop_profile_platforms else None
|
||||
payload["platforms"] = {
|
||||
k: v for k, v in platforms.items() if not isinstance(k, str) or ":" not in k
|
||||
k: v for k, v in platforms.items()
|
||||
if not isinstance(k, str) or ":" not in k or (drop_prefix is not None and not k.startswith(drop_prefix))
|
||||
}
|
||||
# Re-stamp identity + code fields on every write: the file can outlive its creator and the
|
||||
# top-level record must describe the CURRENT writer.
|
||||
|
||||
@@ -887,16 +887,20 @@ async def update_messaging_platform(platform_id: str, body: MessagingPlatformUpd
|
||||
)
|
||||
# A live multiplexer serving this named profile builds the adapter from the new token now
|
||||
# (its periodic rescan would otherwise pick it up within a cycle); no gateway restart.
|
||||
hot_served = bool(target_profile) and await asyncio.to_thread(_notify_multiplexer_hot_serve, target_profile)
|
||||
hot_served = await asyncio.to_thread(_notify_multiplexer_hot_serve, target_profile)
|
||||
return {"ok": True, "platform": platform_id, "hot_served": hot_served}
|
||||
|
||||
|
||||
def _notify_multiplexer_hot_serve(profile: str) -> bool:
|
||||
from hermes_cli.gateway import named_profile_served_by_running_multiplexer
|
||||
def _notify_multiplexer_hot_serve(profile: Optional[str]) -> bool:
|
||||
"""True when a live multiplexer serves the written profile and was told to rebuild its adapters.
|
||||
Unscoped (no ``?profile=``) means THIS process's profile: Desktop routes a pooled
|
||||
``hermes --profile X serve`` without the query (#109088), so X must resolve here too."""
|
||||
from hermes_cli.gateway import _current_profile_name, named_profile_served_by_running_multiplexer
|
||||
from hermes_cli.gateway_multiplex_served import notify_multiplexer_profiles_changed
|
||||
if not named_profile_served_by_running_multiplexer(profile):
|
||||
name = (profile or "").strip() or _current_profile_name()
|
||||
if not name or name == "default" or not named_profile_served_by_running_multiplexer(name):
|
||||
return False
|
||||
return notify_multiplexer_profiles_changed(profile) is not None
|
||||
return notify_multiplexer_profiles_changed(name) is not None
|
||||
|
||||
|
||||
@router.post("/api/messaging/platforms/{platform_id}/test")
|
||||
|
||||
@@ -305,3 +305,37 @@ def test_scoped_enablement_uses_only_own_credentials(client, isolated_profiles,
|
||||
assert platform["enabled"] is enabled
|
||||
assert platform["configured"] is True
|
||||
assert "root-token" in (isolated_profiles["default"] / ".env").read_text(encoding="utf-8")
|
||||
|
||||
|
||||
@pytest.mark.parametrize("topology", ["scoped_query", "pooled_unscoped"])
|
||||
def test_credential_write_hot_serves_a_multiplexed_profile(client, isolated_profiles, monkeypatch, topology):
|
||||
"""A token saved for a profile the live multiplexer serves is handed to the multiplexer right
|
||||
away (``hot_served``), so the UI skips its restart banner. Both Desktop topologies: the dashboard's
|
||||
``?profile=`` and a pooled ``hermes --profile X serve`` that receives the PUT unscoped (#109088)."""
|
||||
import hermes_cli.gateway as gateway_cli
|
||||
import hermes_cli.gateway_multiplex_served as served_mod
|
||||
notified = []
|
||||
monkeypatch.setattr(gateway_cli, "named_profile_served_by_running_multiplexer", lambda name=None: name == "worker_alpha")
|
||||
monkeypatch.setattr(served_mod, "notify_multiplexer_profiles_changed", lambda name, **kw: notified.append(name) or ["default", name])
|
||||
if topology == "pooled_unscoped":
|
||||
monkeypatch.setattr(gateway_cli, "_current_profile_name", lambda: "worker_alpha")
|
||||
params = {}
|
||||
else:
|
||||
params = {"profile": "worker_alpha"}
|
||||
resp = client.put("/api/messaging/platforms/telegram", params=params,
|
||||
json={"enabled": True, "env": {"TELEGRAM_BOT_TOKEN": _VALID_WORKER_BOT_TOKEN}})
|
||||
assert resp.status_code == 200
|
||||
assert resp.json()["hot_served"] is True
|
||||
assert notified == ["worker_alpha"]
|
||||
|
||||
|
||||
def test_credential_write_on_default_profile_is_not_hot_served(client, isolated_profiles, monkeypatch):
|
||||
"""The default profile is the multiplexer itself (its own adapters are restart-managed): never
|
||||
claim a hot serve for it."""
|
||||
import hermes_cli.gateway_multiplex_served as served_mod
|
||||
monkeypatch.setattr(served_mod, "notify_multiplexer_profiles_changed",
|
||||
lambda name, **kw: pytest.fail("default profile must not ping the multiplexer"))
|
||||
resp = client.put("/api/messaging/platforms/telegram",
|
||||
json={"enabled": True, "env": {"TELEGRAM_BOT_TOKEN": _VALID_WORKER_BOT_TOKEN}})
|
||||
assert resp.status_code == 200
|
||||
assert resp.json()["hot_served"] is False
|
||||
|
||||
Reference in New Issue
Block a user