diff --git a/gateway/run.py b/gateway/run.py index 26a3a67f59..a072dac7dd 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -5071,6 +5071,7 @@ async def _start_gateway_start_control_socket(runner): # failure only means consumers fall back to the process-scan/state-file layer, exactly as before # this feature. See #92091. from gateway.control_socket import GatewayControlServer + from gateway.run_profile_reconcile import migrate_profile_identity_verb # pause-for-update: the updater asks us to drain + exit (freeing venv handles) vs. a tree-kill # (same path as SIGUSR1). Handler runs on the socket executor thread, so marshal onto the loop. # pause-for-update (#92091 step 2): the updater asks this gateway to drain in-flight turns and exit @@ -5116,43 +5117,10 @@ async def _start_gateway_start_control_socket(runner): except concurrent.futures.TimeoutError: return {"multiplex": True, "pending": True, "served_profiles": runner.served_profile_names()} - def _migrate_profile_identity_handler(params: dict) -> dict: - """Migrate both durable stores and the routing index owned by this live gateway.""" - old, new = str(params.get("old") or "").strip(), str(params.get("new") or "").strip() - if not old or not new or old == new: - return {"ok": False, "error": "old/new required and must differ"} - store = getattr(runner, "session_store", None) - if store is None: - return {"ok": False, "error": "live gateway has no session store"} - acquired = [] - try: - from hermes_state_registry import acquire, release_or_close - db_counts: dict[str, dict[str, int]] = {} - routing_db = getattr(store, "_routing_db", None) - if routing_db is not None and hasattr(routing_db, "rekey_profile_state"): - db_counts["routing"] = routing_db.rekey_profile_state(old, new) - routing_home = getattr(store, "_routing_home", None) - profile_path = Path(routing_home) / "profiles" / new / "state.db" if routing_home else None - if profile_path is not None and profile_path.exists(): - profile_db = acquire(profile_path) - acquired.append(profile_db) - db_counts["profile"] = profile_db.rekey_profile_state(old, new) - rekeyed = store.rekey_profile_routing(old, new) - return {"ok": True, "rekeyed": rekeyed, "db": db_counts} - except Exception as exc: - logger.warning("Profile identity migration failed for %r->%r: %s", old, new, exc) - return {"ok": False, "error": f"{type(exc).__name__}: {exc}"} - finally: - for db in acquired: - try: - release_or_close(db) - except Exception: - logger.debug("Failed to release renamed profile state DB", exc_info=True) - _control_server = GatewayControlServer( verb_handlers={"pause-for-update": _pause_for_update_handler, "rescan-profiles": _rescan_profiles_handler, - "migrate-profile-identity": _migrate_profile_identity_handler}) + "migrate-profile-identity": migrate_profile_identity_verb(runner)}) if not await _control_server.start(): _control_server = None else: diff --git a/gateway/run_profile_reconcile.py b/gateway/run_profile_reconcile.py index 4aa4ed975f..c127770f2b 100644 --- a/gateway/run_profile_reconcile.py +++ b/gateway/run_profile_reconcile.py @@ -262,3 +262,45 @@ def _mcp_config_reconciler(runner=None): _reconcile_current(str(profile_name)) return _tick + + +def migrate_profile_identity_verb(runner): + """Build the ``migrate-profile-identity`` control-verb handler for ``hermes profile rename`` + (#111926). The live multiplexer owns the routing index in memory and writes it back + periodically, so a CLI-side rewrite of ``agent::*`` would be clobbered on the next save; + the CLI therefore asks this process to rekey both durable stores AND ``SessionStore._entries``. + Runs on the control-socket executor thread; ``rekey_profile_routing`` takes the store lock.""" + + def _handler(params: dict) -> dict: + old, new = str(params.get("old") or "").strip(), str(params.get("new") or "").strip() + if not old or not new or old == new: + return {"ok": False, "error": "old/new required and must differ"} + store = getattr(runner, "session_store", None) + if store is None: + return {"ok": False, "error": "live gateway has no session store"} + acquired = [] + try: + from hermes_state_registry import acquire, release_or_close + db_counts: Dict[str, Dict[str, int]] = {} + routing_db = getattr(store, "_routing_db", None) + if routing_db is not None and hasattr(routing_db, "rekey_profile_state"): + db_counts["routing"] = routing_db.rekey_profile_state(old, new) + routing_home = getattr(store, "_routing_home", None) + profile_path = Path(routing_home) / "profiles" / new / "state.db" if routing_home else None + if profile_path is not None and profile_path.exists(): + profile_db = acquire(profile_path) + acquired.append(profile_db) + db_counts["profile"] = profile_db.rekey_profile_state(old, new) + rekeyed = store.rekey_profile_routing(old, new) + return {"ok": True, "rekeyed": rekeyed, "db": db_counts} + except Exception as exc: + logger.warning("Profile identity migration failed for %r->%r: %s", old, new, exc) + return {"ok": False, "error": f"{type(exc).__name__}: {exc}"} + finally: + for db in acquired: + try: + release_or_close(db) + except Exception: + logger.debug("Failed to release renamed profile state DB", exc_info=True) + + return _handler diff --git a/hermes_cli/profile_cmd.py b/hermes_cli/profile_cmd.py index 2d355fe721..f20165a35f 100644 --- a/hermes_cli/profile_cmd.py +++ b/hermes_cli/profile_cmd.py @@ -440,7 +440,7 @@ def _profile_migrate_identity(args): """Retry the identity migration of a rename that already completed. Exits non-zero when a live gateway would not migrate (it still owns the routing index in memory), or when a database rejected the rewrite (collision, lock, partial failure).""" - from hermes_cli.profiles import migrate_profile_identity + from hermes_cli.profile_identity import migrate_profile_identity try: migrated = migrate_profile_identity(args.old_name, args.new_name) except (ValueError, FileNotFoundError) as e: diff --git a/hermes_cli/profile_identity.py b/hermes_cli/profile_identity.py new file mode 100644 index 0000000000..b70e81ee17 --- /dev/null +++ b/hermes_cli/profile_identity.py @@ -0,0 +1,119 @@ +"""Rekey a renamed profile's session/routing identity (#111926). + +``rename_profile`` moves ``profiles//`` to ``profiles//`` so row DATA travels with the +directory, but the profile name is also baked into keys and values the move never touches: +``agent::*`` session-key namespaces (routing index + the profile's own ``sessions`` rows), +``sessions.profile_name``, ``gateway_heartbeats.profile`` and ``delivery_obligations``. Left alone, +every inbound event on a chat keyed to the old name logs ``Profile 'old' does not exist`` and +falls back to the global home, and renamed sessions drop out of the Desktop sidebar. + +Ownership decides who rewrites: a live multiplexer holds the routing index in memory +(``SessionStore._entries``) and writes it back periodically, so a CLI-side DB rewrite would be +clobbered on its next save — the CLI delegates to the ``migrate-profile-identity`` control verb. +With no live multiplexer nothing else holds the store and the durable rewrite is safe here. +""" +from __future__ import annotations + +import contextlib +import sys +from pathlib import Path + + +def migrate_profile_identity(old_name: str, new_name: str) -> bool: + """Retry the session/routing identity migration of a rename that already completed. + + ``rename_profile`` runs the migration itself; this is the standalone retry behind + ``hermes profile migrate-identity `` for when that attempt failed. The rename + cannot simply be repeated — ``profiles/`` is gone — and the identity to migrate is read + from the DB rows that still name *old*, so only the new profile has to exist here. + + A live multiplexer holds the routing index in memory and therefore stays the owner of the + migration (the CLI delegates to its control verb); with no live multiplexer the durable + rewrite is safe because nothing else holds the store. Idempotent: re-running a completed + migration succeeds with nothing left to rekey. Returns True when the identity was migrated, + False when a live gateway would not do it — the caller reports that as a failure. + """ + from hermes_cli.profiles import _canon_valid, _live_default_multiplexer, _unknown_profile_error, get_profile_dir + old_canon = _canon_valid(old_name) + new_canon = _canon_valid(new_name) + if "default" in (old_canon, new_canon): + raise ValueError("Identity migration applies to named profiles only.") + if not get_profile_dir(new_canon).is_dir(): + raise _unknown_profile_error(new_canon) + return _migrate_profile_identity(old_canon, new_canon, _live_default_multiplexer()) + + +def _control_answer_failure(answer) -> str: + """Why a control-socket answer is not a success. Keeps the raw answer when the payload carries + no reason field, so a malformed or old-gateway response stays diagnosable instead of + collapsing into a generic warning.""" + if isinstance(answer, dict): + failure = answer.get("error") or answer.get("message") or answer.get("detail") + return str(failure) if failure else repr(answer) + if answer is not None: + return repr(answer) + return "no response from gateway control socket" + + +def _gateway_accepts_profile_identity_verb(root: Path) -> bool: + """True when the gateway at *root* answers a verb it has always had. Distinguishes a failed + migration verb caused by an older gateway process from one caused by no gateway at all.""" + try: + from gateway.control_socket import identify_gateway + return identify_gateway(root) is not None + except Exception: + return False + + +def _migrate_profile_identity(old_canon: str, new_canon: str, live_mux: bool) -> bool: + """Rekey renamed-profile identity without racing a live gateway's in-memory routing index. + + Returns True when the identity was migrated — by the gateway's control verb, or by this + process's durable rewrite when no gateway holds the store — and False when a live gateway did + not accept it. Never fatal to the rename, which has already happened by this point. + """ + if live_mux: + from hermes_constants import get_default_hermes_root + root = get_default_hermes_root() + try: + from gateway.control_socket import migrate_gateway_profile_identity + answer = migrate_gateway_profile_identity(root, old_canon, new_canon) + except Exception as exc: + reason = f"{type(exc).__name__}: {exc}" + else: + if isinstance(answer, dict) and answer.get("ok") is True: + return True + reason = _control_answer_failure(answer) + if answer is None and _gateway_accepts_profile_identity_verb(root): + reason += (" — the gateway is running but does not implement " + "'migrate-profile-identity' (an older process than this CLI)") + print( + "⚠ Profile was renamed, but the live gateway could not migrate session identity" + f" ({reason}). Restart the gateway, then run:\n" + f" hermes profile migrate-identity {old_canon} {new_canon}", + file=sys.stderr) + return False + + from hermes_cli.profiles import get_profile_dir + from hermes_state_registry import acquire, release_or_close + from hermes_constants import get_default_hermes_root + root = get_default_hermes_root() + migrated = True + for db_path in (root / "state.db", get_profile_dir(new_canon) / "state.db"): + if not db_path.exists(): + continue + db = None + try: + db = acquire(db_path) + db.rekey_profile_state(old_canon, new_canon) + except Exception as exc: + migrated = False + print( + f"⚠ Profile was renamed, but identity migration failed for {db_path}: " + f"{type(exc).__name__}: {exc}", + file=sys.stderr) + finally: + if db is not None: + with contextlib.suppress(Exception): + release_or_close(db) + return migrated diff --git a/hermes_cli/profiles.py b/hermes_cli/profiles.py index 48fbe6ae2a..a3ce6412c4 100644 --- a/hermes_cli/profiles.py +++ b/hermes_cli/profiles.py @@ -1867,6 +1867,7 @@ def rename_profile(old_name: str, new_name: str) -> Path: # 6. Migrate profile-name-keyed session/routing state (session keys, profile_name, heartbeats, # delivery + routing index) from the old name to the new one. A stale ``agent::*`` routing # key otherwise resolves to a profile that no longer exists on every inbound event. + from hermes_cli.profile_identity import _migrate_profile_identity _migrate_profile_identity(old_canon, new_canon, live_mux) # 7. Hot-serve the renamed profile now (mirrors create; a missed signal only delays it). @@ -1875,104 +1876,6 @@ def rename_profile(old_name: str, new_name: str) -> Path: return new_dir -def migrate_profile_identity(old_name: str, new_name: str) -> bool: - """Retry the session/routing identity migration of a rename that already completed. - - ``rename_profile`` runs the migration itself; this is the standalone retry behind - ``hermes profile migrate-identity `` for when that attempt failed. The rename - cannot simply be repeated — ``profiles/`` is gone — and the identity to migrate is read - from the DB rows that still name *old*, so only the new profile has to exist here. - - A live multiplexer holds the routing index in memory and therefore stays the owner of the - migration (the CLI delegates to its control verb); with no live multiplexer the durable - rewrite is safe because nothing else holds the store. Idempotent: re-running a completed - migration succeeds with nothing left to rekey. Returns True when the identity was migrated, - False when a live gateway would not do it — the caller reports that as a failure. - """ - old_canon = _canon_valid(old_name) - new_canon = _canon_valid(new_name) - if "default" in (old_canon, new_canon): - raise ValueError("Identity migration applies to named profiles only.") - if not get_profile_dir(new_canon).is_dir(): - raise _unknown_profile_error(new_canon) - return _migrate_profile_identity(old_canon, new_canon, _live_default_multiplexer()) - - -def _control_answer_failure(answer) -> str: - """Why a control-socket answer is not a success. Keeps the raw answer when the payload carries - no reason field, so a malformed or old-gateway response stays diagnosable instead of - collapsing into a generic warning.""" - if isinstance(answer, dict): - failure = answer.get("error") or answer.get("message") or answer.get("detail") - return str(failure) if failure else repr(answer) - if answer is not None: - return repr(answer) - return "no response from gateway control socket" - - -def _gateway_accepts_profile_identity_verb(root: Path) -> bool: - """True when the gateway at *root* answers a verb it has always had. Distinguishes a failed - migration verb caused by an older gateway process from one caused by no gateway at all.""" - try: - from gateway.control_socket import identify_gateway - return identify_gateway(root) is not None - except Exception: - return False - - -def _migrate_profile_identity(old_canon: str, new_canon: str, live_mux: bool) -> bool: - """Rekey renamed-profile identity without racing a live gateway's in-memory routing index. - - Returns True when the identity was migrated — by the gateway's control verb, or by this - process's durable rewrite when no gateway holds the store — and False when a live gateway did - not accept it. Never fatal to the rename, which has already happened by this point. - """ - if live_mux: - from hermes_constants import get_default_hermes_root - root = get_default_hermes_root() - try: - from gateway.control_socket import migrate_gateway_profile_identity - answer = migrate_gateway_profile_identity(root, old_canon, new_canon) - except Exception as exc: - reason = f"{type(exc).__name__}: {exc}" - else: - if isinstance(answer, dict) and answer.get("ok") is True: - return True - reason = _control_answer_failure(answer) - if answer is None and _gateway_accepts_profile_identity_verb(root): - reason += (" — the gateway is running but does not implement " - "'migrate-profile-identity' (an older process than this CLI)") - print( - "⚠ Profile was renamed, but the live gateway could not migrate session identity" - f" ({reason}). Restart the gateway, then run:\n" - f" hermes profile migrate-identity {old_canon} {new_canon}", - file=sys.stderr) - return False - - from hermes_state_registry import acquire, release_or_close - from hermes_constants import get_default_hermes_root - root = get_default_hermes_root() - migrated = True - for db_path in (root / "state.db", get_profile_dir(new_canon) / "state.db"): - if not db_path.exists(): - continue - db = None - try: - db = acquire(db_path) - db.rekey_profile_state(old_canon, new_canon) - except Exception as exc: - migrated = False - print( - f"⚠ Profile was renamed, but identity migration failed for {db_path}: " - f"{type(exc).__name__}: {exc}", - file=sys.stderr) - finally: - if db is not None: - with contextlib.suppress(Exception): - release_or_close(db) - return migrated - - # Profile env resolution (called from _apply_profile_override) def resolve_profile_env(profile_name: str) -> str: diff --git a/tests/gateway/test_rekey_profile_routing.py b/tests/gateway/test_rekey_profile_routing.py index daef3d6d97..f57944bb7b 100644 --- a/tests/gateway/test_rekey_profile_routing.py +++ b/tests/gateway/test_rekey_profile_routing.py @@ -51,16 +51,6 @@ def test_rekeys_old_namespace_and_origin_profile(tmp_path): assert store._entries["agent:keepme:feishu:dm:chatB"].origin.profile == "keepme" -def test_noop_for_equal_or_empty_names(tmp_path): - store = _make_store(tmp_path) - with store._lock: - store._entries["agent:oldname:feishu:dm:chatA"] = _entry( - "agent:oldname:feishu:dm:chatA", "chatA", "oldname") - assert store.rekey_profile_routing("x", "x") == 0 - assert store.rekey_profile_routing("", "y") == 0 - assert "agent:oldname:feishu:dm:chatA" in store._entries - - def test_does_not_overwrite_existing_new_namespace_key(tmp_path): store = _make_store(tmp_path) with store._lock: diff --git a/tests/hermes_cli/test_profiles.py b/tests/hermes_cli/test_profiles.py index 4a5282a506..b2ba3f30f3 100644 --- a/tests/hermes_cli/test_profiles.py +++ b/tests/hermes_cli/test_profiles.py @@ -993,33 +993,6 @@ class TestRenameProfile: assert "agent:newname:feishu:dm:chatA" in routing root_db2.close() - def test_migrate_identity_reports_the_raw_answer_and_exits_nonzero(self, profile_env, capsys): - """A live gateway that answers with something unusable must fail loudly — exit non-zero, - name the retry command, and quote the raw answer (a non-dict payload used to print a - reason-less warning). The live gateway keeps ownership: no direct DB rewrite.""" - from hermes_cli.profile_cmd import cmd_profile - from argparse import Namespace - create_profile("oldname", no_alias=True) - create_profile("newname", no_alias=True) # the rename already happened; only must exist - - with patch("hermes_cli.profiles._live_default_multiplexer", return_value=True), \ - patch("gateway.control_socket.migrate_gateway_profile_identity", - return_value="not a control answer"), \ - patch("hermes_state_registry.acquire") as acquire: - with pytest.raises(SystemExit) as excinfo: - cmd_profile(Namespace(profile_action="migrate-identity", - old_name="oldname", new_name="newname")) - - assert excinfo.value.code != 0 - err = capsys.readouterr().err - assert "not a control answer" in err - assert "hermes profile migrate-identity oldname newname" in err - acquire.assert_not_called() - - -# =================================================================== -# TestExportImport -# =================================================================== class TestExportImport: """Tests for export_profile() / import_profile().""" diff --git a/tests/hermes_state/test_rekey_profile_state.py b/tests/hermes_state/test_rekey_profile_state.py index 307338fbd2..92ba91a1fc 100644 --- a/tests/hermes_state/test_rekey_profile_state.py +++ b/tests/hermes_state/test_rekey_profile_state.py @@ -100,23 +100,6 @@ class TestRekeyProfileState: assert ob["session_key"] == "agent:newname:feishu:dm:chatA" assert ob["adapter_profile"] == "newname" - def test_noop_when_names_equal_or_empty(self, db): - assert db.rekey_profile_state("x", "x") == {} - assert db.rekey_profile_state("", "y") == {} - assert db.rekey_profile_state("x", "") == {} - - def test_idempotent(self, db): - db.create_session( - "sess_old", "feishu", session_key="agent:oldname:feishu:dm:chatA", - profile_name="oldname", chat_id="chatA", chat_type="dm", - ) - first = db.rekey_profile_state("oldname", "newname") - assert first["sessions_session_key"] == 1 - second = db.rekey_profile_state("oldname", "newname") - # Nothing left under the old name. - assert second["sessions_session_key"] == 0 - assert second["sessions_profile_name"] == 0 - def test_underscore_in_profile_name_is_not_a_like_wildcard(self, db): db.create_session( "literal", "telegram", session_key="agent:foo_bar:telegram:dm:a", @@ -154,4 +137,3 @@ class TestRekeyProfileState: "WHERE chat_id = ? AND thread_id = ?", ("chatA", "threadA")) assert binding["profile_name"] == "newname" assert binding["session_key"] == "agent:newname:telegram:dm:chatA" - diff --git a/tests/utils/test_atomic_writers_deleted_profile.py b/tests/test_utils_atomic_writers_deleted_profile.py similarity index 77% rename from tests/utils/test_atomic_writers_deleted_profile.py rename to tests/test_utils_atomic_writers_deleted_profile.py index 95daee405d..7da2d29e8b 100644 --- a/tests/utils/test_atomic_writers_deleted_profile.py +++ b/tests/test_utils_atomic_writers_deleted_profile.py @@ -50,30 +50,6 @@ class TestAtomicWritersRefuseDeletedProfileHome: atomic_json_write(profile / "cache" / "reasoning_caps.json", {"m": {}}) assert not profile.exists() - def test_atomic_write_text_does_not_recreate_home(self, tmp_path): - profile = _tombstoned_profile(tmp_path) - with pytest.raises( - FileNotFoundError, match="Named profile home does not exist" - ): - atomic_write_text(profile / "cache" / "models-dev-etag", "etag") - assert not profile.exists() - - def test_late_reasoning_caps_save_after_delete(self, tmp_path): - from hermes_cli import models_reasoning_caps - - profile = _tombstoned_profile(tmp_path) - token = set_hermes_home_override(profile) - try: - models_reasoning_caps._save_reasoning_caps_disk( - "https://example/v1/models", {"m": {"supports_reasoning": True}} - ) - finally: - from hermes_constants import reset_hermes_home_override - - reset_hermes_home_override(token) - assert not (profile / "cache" / "reasoning_caps.json").exists() - assert not profile.exists() - def test_late_models_cache_save_after_delete(self, tmp_path): from hermes_cli.models import _write_json_cache @@ -147,12 +123,3 @@ class TestUnrelatedProfilesPathsStillWrite: assert json.loads( (custom_home / "cache" / "blob.json").read_text(encoding="utf-8") ) == {"a": 1} - - def test_default_home_cache_write(self, tmp_path): - home = tmp_path / ".hermes" - atomic_write_text(home / "cache" / "etag", "v1") - assert (home / "cache" / "etag").read_text(encoding="utf-8") == "v1" - - def test_plain_tmp_path_write(self, tmp_path): - atomic_json_write(tmp_path / "plain" / "data.json", [1, 2]) - assert (tmp_path / "plain" / "data.json").exists()