From eb4f92bc23a346db41d4f65582e1ca077d0bc89f Mon Sep 17 00:00:00 2001 From: JoaoMarcos44 Date: Thu, 24 Sep 2026 01:51:51 -0300 Subject: [PATCH 01/22] fix(dashboard): reject credential preview writeback (cherry picked from commit 9c8bc79ade733ab853d7cc33aa1021939b16b363) --- hermes_cli/web_routers/config_env.py | 19 ++++++++++++++++--- hermes_cli/web_routers/messaging.py | 27 +++++++++++++++++++++------ 2 files changed, 37 insertions(+), 9 deletions(-) diff --git a/hermes_cli/web_routers/config_env.py b/hermes_cli/web_routers/config_env.py index 606495eca3..28e092ecde 100644 --- a/hermes_cli/web_routers/config_env.py +++ b/hermes_cli/web_routers/config_env.py @@ -47,6 +47,10 @@ _reveal_timestamps: List[float] = [] _REVEAL_MAX_PER_WINDOW = 5 _REVEAL_WINDOW_SECONDS = 30 +_REDACTED_CREDENTIAL_WRITE_DETAIL = ( + "Refusing to save a redacted credential preview; re-enter the full secret to replace it." +) + # Display order for tabs — unlisted categories sort alphabetically after these. _CATEGORY_ORDER = [ "general", "agent", "terminal", "display", "delegation", @@ -295,9 +299,13 @@ async def set_env_var(body: EnvVarUpdate, profile: Optional[str] = None): with _env_write_errors("PUT /api/env failed", http_passthrough=False): from hermes_cli.credential_lifecycle import save_provider_env_credential - return await scoped_to_thread( - body.profile or profile, lambda: save_provider_env_credential(body.key, body.value) - ) + def _save(): + current = load_env().get(body.key) + if current and body.value == redact_key(current): + raise ValueError(_REDACTED_CREDENTIAL_WRITE_DETAIL) + return save_provider_env_credential(body.key, body.value) + + return await scoped_to_thread(body.profile or profile, _save) # Live credential probes keyed by env var: (url, auth) where auth is "bearer" @@ -620,6 +628,11 @@ def _write_custom_endpoint(cfg: Dict[str, Any], body: CustomEndpointUpdate) -> T env_var = custom_endpoint_key_env(endpoint_id) submitted_key = body.api_key.strip() if body.api_key is not None else None if submitted_key: + display_preview = _api_key_display(existing)[1] + existing_key_env = str(existing.get("key_env") or "").strip() + stored_secret = load_env().get(existing_key_env) if existing_key_env else None + if submitted_key == display_preview or (stored_secret and submitted_key == redact_key(stored_secret)): + raise HTTPException(status_code=400, detail=_REDACTED_CREDENTIAL_WRITE_DETAIL) save_env_value(env_var, submitted_key) entry["key_env"] = env_var entry.pop("api_key", None) diff --git a/hermes_cli/web_routers/messaging.py b/hermes_cli/web_routers/messaging.py index 5d96faf743..7dde2c7cc7 100644 --- a/hermes_cli/web_routers/messaging.py +++ b/hermes_cli/web_routers/messaging.py @@ -871,17 +871,32 @@ async def update_messaging_platform(platform_id: str, body: MessagingPlatformUpd raise HTTPException(status_code=400, detail=f"{key} is not configurable for {entry['name']}") def _apply(): - with _profile_scope(target_profile): + with _profile_scope(target_profile) as scoped_dir: + env_on_disk = load_env() + updates: dict[str, str] = {} + + # Validate the whole request against the same credential snapshot the GET + # response uses before clearing or replacing anything. for key in body.clear_env: _check_allowed(key) - remove_env_value(key) - for key, value in body.env.items(): _check_allowed(key) trimmed = value.strip() - if trimmed: - _validate_messaging_env_value(platform_id, key, trimmed) - save_env_value(key, trimmed) + if not trimmed: + continue + current = env_on_disk.get(key) or ("" if scoped_dir is not None else os.getenv(key, "")) + if current and trimmed == redact_key(current): + raise HTTPException( + status_code=400, + detail="Refusing to save a redacted credential preview; re-enter the full secret to replace it.", + ) + _validate_messaging_env_value(platform_id, key, trimmed) + updates[key] = trimmed + + for key in body.clear_env: + remove_env_value(key) + for key, value in updates.items(): + save_env_value(key, value) if body.enabled is not None: _write_platform_enabled(platform_id, body.enabled) From e3ef1418cebd96558c76b37026007832489e41ec Mon Sep 17 00:00:00 2001 From: JoaoMarcos44 Date: Thu, 24 Sep 2026 01:52:56 -0300 Subject: [PATCH 02/22] test(dashboard): cover secret preview boundaries (cherry picked from commit 10381effd3a2bc74c1c67e1b6deb3bd822655564) --- tests/hermes_cli/test_web_server.py | 63 +++++++++++++++++++++++++++++ 1 file changed, 63 insertions(+) diff --git a/tests/hermes_cli/test_web_server.py b/tests/hermes_cli/test_web_server.py index 790590e031..92fff0b1f0 100644 --- a/tests/hermes_cli/test_web_server.py +++ b/tests/hermes_cli/test_web_server.py @@ -2142,6 +2142,69 @@ CONFIG_SCHEMA = ProviderConfigSchema( assert endpoint["has_api_key"] is True assert "sk-in-env" not in (endpoint["api_key_preview"] or "") + def test_env_rejects_its_redacted_preview(self): + from hermes_cli.config import load_env, save_env_value + + real = "sk-live-secret-abcdef1234567890" + save_env_value("OPENAI_API_KEY", real) + preview = self.client.get("/api/env").json()["OPENAI_API_KEY"]["redacted_value"] + + response = self.client.put("/api/env", json={"key": "OPENAI_API_KEY", "value": preview}) + + assert response.status_code == 400 + assert load_env()["OPENAI_API_KEY"] == real + + def test_custom_endpoint_rejects_current_and_legacy_previews(self): + from hermes_cli.config import custom_endpoint_key_env, get_env_value + + payload = { + "id": "proxy", "name": "Proxy", "base_url": "https://llm.example.com/v1", + "model": "m", "api_key": "sk-live-secret-abcdef1234567890", + } + created = self.client.post("/api/providers/custom-endpoints", json=payload) + endpoint = next(e for e in created.json()["endpoints"] if e["id"] == "proxy") + env_key = custom_endpoint_key_env("proxy") + real = get_env_value(env_key) + + for preview in (endpoint["api_key_preview"], redact_key(real)): + response = self.client.post( + "/api/providers/custom-endpoints", json={**payload, "api_key": preview}) + assert response.status_code == 400 + assert get_env_value(env_key) == real + + def test_messaging_rejects_preview_from_process_env(self, monkeypatch): + from hermes_cli.config import load_env + + key = "DISCORD_BOT_TOKEN" + real = "discord-live-secret-abcdef1234567890" + monkeypatch.setenv(key, real) + assert key not in load_env() + + platforms = self.client.get("/api/messaging/platforms").json()["platforms"] + discord = next(p for p in platforms if p["id"] == "discord") + preview = next(v["redacted_value"] for v in discord["env_vars"] if v["key"] == key) + response = self.client.put(f"/api/messaging/platforms/discord", json={"env": {key: preview}}) + + assert response.status_code == 400 + assert key not in load_env() + assert os.environ[key] == real + + def test_messaging_rejects_preview_before_clearing_secret(self): + from hermes_cli.config import load_env, save_env_value + + key = "DISCORD_BOT_TOKEN" + real = "discord-live-secret-abcdef1234567890" + save_env_value(key, real) + preview = redact_key(real) + + response = self.client.put( + "/api/messaging/platforms/discord", + json={"clear_env": [key], "env": {key: preview}}, + ) + + assert response.status_code == 400 + assert load_env()[key] == real + def test_activating_an_endpoint_carries_its_credential_either_way(self): """Activate must work for both key_env and pre-#69449 plaintext entries.""" from hermes_cli.config import load_config, save_config From f8b66c75318ab5ad84a5b904b8488000bc4a5155 Mon Sep 17 00:00:00 2001 From: JoaoMarcos44 Date: Thu, 24 Sep 2026 02:11:47 -0300 Subject: [PATCH 03/22] fix(dashboard): make credential previews non-reusable (cherry picked from commit 215b51316aa2c1037223c4807ef6ec9b631fb7ab) --- hermes_cli/web_routers/_common.py | 25 ++++++ hermes_cli/web_routers/config_env.py | 27 ++++--- hermes_cli/web_routers/messaging.py | 27 ++++--- tests/hermes_cli/test_web_server.py | 112 +++++++++++++++++++++++++-- 4 files changed, 160 insertions(+), 31 deletions(-) diff --git a/hermes_cli/web_routers/_common.py b/hermes_cli/web_routers/_common.py index 2399f90e24..67b44b2556 100644 --- a/hermes_cli/web_routers/_common.py +++ b/hermes_cli/web_routers/_common.py @@ -118,6 +118,31 @@ def require(value: Optional[str], detail: str) -> str: return stripped +REDACTED_CREDENTIAL_WRITE_DETAIL = ( + "Refusing to save a redacted credential preview; re-enter the full secret to replace it." +) + + +def redacted_credential_preview(value: Any) -> Optional[str]: + """Return a display-only credential sentinel that can never gain write authority.""" + if not value: + return None + from hermes_cli.config import redact_key + return f"«redacted:{redact_key(str(value))}»" + + +def is_redacted_credential_preview(submitted: Any, current: Any = None) -> bool: + """Recognize current and stale dashboard previews without guessing from key shape.""" + value = str(submitted or "") + if value == "«redacted-secret»" or (value.startswith("«redacted:") and value.endswith("»")): + return True + if current: + # Compatibility with a page opened before the non-reusable sentinel contract. + from hermes_cli.config import redact_key + return value == redact_key(str(current)) + return False + + # Corrupt-store reporting for polled read endpoints. The dashboard polls analytics every few # seconds; a persistently malformed state.db once produced ~520K identical tracebacks in 24 h # (#96591). One WARNING per store per interval, then debug; the caller gets an explicit status diff --git a/hermes_cli/web_routers/config_env.py b/hermes_cli/web_routers/config_env.py index 28e092ecde..0487c81147 100644 --- a/hermes_cli/web_routers/config_env.py +++ b/hermes_cli/web_routers/config_env.py @@ -11,7 +11,10 @@ import asyncio import time import urllib.parse from fastapi import APIRouter -from hermes_cli.web_routers._common import http_failure, scoped_to_thread +from hermes_cli.web_routers._common import ( + REDACTED_CREDENTIAL_WRITE_DETAIL, http_failure, is_redacted_credential_preview, + redacted_credential_preview, scoped_to_thread, +) from hermes_cli.web_deps import LateState, late from hermes_cli.web_server_config import ( _apply_main_model_assignment, _denormalize_config_from_web, _normalize_config_for_web, _schema_with_dynamic_provider_options, @@ -21,7 +24,7 @@ from hermes_cli.web_server_profiles import ( _approval_mode_of, _broadcast_gateway_session_info, _is_other_profile, _parse_model_entries, ) from fastapi import HTTPException, Request -from hermes_cli.config import DEFAULT_CONFIG, OPTIONAL_ENV_VARS, read_raw_config, require_readable_config_before_write, custom_endpoint_key_env, coerce_provider_id, find_provider_entry, get_compatible_custom_providers, redact_key, _deep_merge +from hermes_cli.config import DEFAULT_CONFIG, OPTIONAL_ENV_VARS, read_raw_config, require_readable_config_before_write, custom_endpoint_key_env, coerce_provider_id, find_provider_entry, get_compatible_custom_providers, _deep_merge from hermes_cli.config_providers import _canonical_api_mode, _custom_provider_entry_to_provider_config from hermes_cli.web_models import ConfigUpdate, EnvVarUpdate, EnvVarDelete, EnvVarReveal, CustomEndpointUpdate from typing import Any, Dict, List, Optional, Tuple @@ -47,10 +50,6 @@ _reveal_timestamps: List[float] = [] _REVEAL_MAX_PER_WINDOW = 5 _REVEAL_WINDOW_SECONDS = 30 -_REDACTED_CREDENTIAL_WRITE_DETAIL = ( - "Refusing to save a redacted credential preview; re-enter the full secret to replace it." -) - # Display order for tabs — unlisted categories sort alphabetically after these. _CATEGORY_ORDER = [ "general", "agent", "terminal", "display", "delegation", @@ -250,7 +249,7 @@ def _get_env_vars_sync(profile: Optional[str] = None): # gaps (description/url) and always supplies provider grouping hints. return { "is_set": bool(value), - "redacted_value": redact_key(value) if value else None, + "redacted_value": redacted_credential_preview(value), "description": info.get("description") or cat_meta.get("description", ""), "url": info.get("url") if info.get("url") is not None else cat_meta.get("url"), "category": info.get("category") or cat_meta.get("category", ""), @@ -301,8 +300,8 @@ async def set_env_var(body: EnvVarUpdate, profile: Optional[str] = None): def _save(): current = load_env().get(body.key) - if current and body.value == redact_key(current): - raise ValueError(_REDACTED_CREDENTIAL_WRITE_DETAIL) + if is_redacted_credential_preview(body.value, current): + raise ValueError(REDACTED_CREDENTIAL_WRITE_DETAIL) return save_provider_env_credential(body.key, body.value) return await scoped_to_thread(body.profile or profile, _save) @@ -368,7 +367,7 @@ def _api_key_display(entry: Dict[str, Any]) -> Tuple[bool, Optional[str]]: """ plaintext = str(entry.get("api_key") or "").strip() if plaintext: - return True, redact_key(plaintext) + return True, redacted_credential_preview(plaintext) key_env = str(entry.get("key_env") or "").strip() if key_env: return True, f"${{{key_env}}}" @@ -631,8 +630,12 @@ def _write_custom_endpoint(cfg: Dict[str, Any], body: CustomEndpointUpdate) -> T display_preview = _api_key_display(existing)[1] existing_key_env = str(existing.get("key_env") or "").strip() stored_secret = load_env().get(existing_key_env) if existing_key_env else None - if submitted_key == display_preview or (stored_secret and submitted_key == redact_key(stored_secret)): - raise HTTPException(status_code=400, detail=_REDACTED_CREDENTIAL_WRITE_DETAIL) + if ( + submitted_key == display_preview + or re.fullmatch(r"\$\{[^}]+\}", submitted_key) + or is_redacted_credential_preview(submitted_key, stored_secret) + ): + raise HTTPException(status_code=400, detail=REDACTED_CREDENTIAL_WRITE_DETAIL) save_env_value(env_var, submitted_key) entry["key_env"] = env_var entry.pop("api_key", None) diff --git a/hermes_cli/web_routers/messaging.py b/hermes_cli/web_routers/messaging.py index 7dde2c7cc7..41f146d992 100644 --- a/hermes_cli/web_routers/messaging.py +++ b/hermes_cli/web_routers/messaging.py @@ -25,14 +25,17 @@ from gateway.status import ( multiplexer_liveness_for_profile, profile_platforms_from_multiplexer, resolve_gateway_liveness, retained_gateway_state) from hermes_cli._subprocess_compat import windows_hide_flags -from hermes_cli.config import OPTIONAL_ENV_VARS, get_env_path, redact_key +from hermes_cli.config import OPTIONAL_ENV_VARS, get_env_path from hermes_constants import get_process_hermes_home from hermes_cli.web_deps import LateState, late from hermes_cli.web_server_gateway import _restart_gateway_after from hermes_cli.web_server_messaging import ( _TelegramOnboardingPairing, _WhatsAppOnboardingSession, _messaging_platform_catalog, _telegram_onboarding_error_message, _telegram_onboarding_lock, _telegram_onboarding_pairings, _whatsapp_onboarding_payload, _whatsapp_onboarding_sessions, ) -from hermes_cli.web_routers._common import http_failure +from hermes_cli.web_routers._common import ( + REDACTED_CREDENTIAL_WRITE_DETAIL, http_failure, is_redacted_credential_preview, + redacted_credential_preview, +) from hermes_cli.web_models import ( MessagingPlatformUpdate, TelegramOnboardingApply, TelegramOnboardingStart, WhatsAppOnboardingApply, WhatsAppOnboardingStart, @@ -196,6 +199,11 @@ def _platform_enablement( return enabled, configured, home_channel +def _messaging_env_value(key: str, env_on_disk: dict[str, str], *, scoped: bool) -> str: + """Resolve the credential source shared by Messaging GET and write preflight.""" + return env_on_disk.get(key) or ("" if scoped else os.getenv(key, "")) + + def _messaging_platform_payload( entry: dict[str, Any], env_on_disk: dict[str, str], runtime: dict | None, scoped: bool = False, profile_home: Optional[Path] = None, @@ -225,14 +233,12 @@ def _messaging_platform_payload( runtime_platform = {} def env_value(key: str) -> str: - # Profile-scoped: judge only the profile's own .env — the dashboard process's - # os.environ carries the ROOT install's .env and would report root credentials as the profile's. - return env_on_disk.get(key) or ("" if scoped else os.getenv(key, "")) + return _messaging_env_value(key, env_on_disk, scoped=scoped) env_vars = [ { "key": key, "required": key in entry["required_env"], "is_set": bool(value), - "redacted_value": redact_key(value) if value else None, **_messaging_env_info(key), + "redacted_value": redacted_credential_preview(value), **_messaging_env_info(key), } for key, value in ((key, env_value(key)) for key in entry["env_vars"]) ] @@ -884,12 +890,9 @@ async def update_messaging_platform(platform_id: str, body: MessagingPlatformUpd trimmed = value.strip() if not trimmed: continue - current = env_on_disk.get(key) or ("" if scoped_dir is not None else os.getenv(key, "")) - if current and trimmed == redact_key(current): - raise HTTPException( - status_code=400, - detail="Refusing to save a redacted credential preview; re-enter the full secret to replace it.", - ) + current = _messaging_env_value(key, env_on_disk, scoped=scoped_dir is not None) + if is_redacted_credential_preview(trimmed, current): + raise HTTPException(status_code=400, detail=REDACTED_CREDENTIAL_WRITE_DETAIL) _validate_messaging_env_value(platform_id, key, trimmed) updates[key] = trimmed diff --git a/tests/hermes_cli/test_web_server.py b/tests/hermes_cli/test_web_server.py index 92fff0b1f0..d38c9da700 100644 --- a/tests/hermes_cli/test_web_server.py +++ b/tests/hermes_cli/test_web_server.py @@ -2145,14 +2145,23 @@ CONFIG_SCHEMA = ProviderConfigSchema( def test_env_rejects_its_redacted_preview(self): from hermes_cli.config import load_env, save_env_value + key = "OPENAI_API_KEY" real = "sk-live-secret-abcdef1234567890" - save_env_value("OPENAI_API_KEY", real) - preview = self.client.get("/api/env").json()["OPENAI_API_KEY"]["redacted_value"] - - response = self.client.put("/api/env", json={"key": "OPENAI_API_KEY", "value": preview}) + save_env_value(key, real) + preview = self.client.get("/api/env").json()[key]["redacted_value"] + assert preview.startswith("«redacted") + response = self.client.put("/api/env", json={"key": key, "value": preview}) assert response.status_code == 400 - assert load_env()["OPENAI_API_KEY"] == real + assert load_env()[key] == real + + # A preview remains display-only even if another actor rotates the secret + # after the GET that minted it. + rotated = "sk-rotated-secret-0987654321" + save_env_value(key, rotated) + response = self.client.put("/api/env", json={"key": key, "value": preview}) + assert response.status_code == 400 + assert load_env()[key] == rotated def test_custom_endpoint_rejects_current_and_legacy_previews(self): from hermes_cli.config import custom_endpoint_key_env, get_env_value @@ -2172,6 +2181,91 @@ CONFIG_SCHEMA = ProviderConfigSchema( assert response.status_code == 400 assert get_env_value(env_key) == real + def test_custom_endpoint_rejects_stale_key_env_placeholder_after_rotation(self): + from hermes_cli.config import load_config, save_config, save_env_value + + provider_id = "env-preview" + old_env = "OLD_ENDPOINT_KEY" + new_env = "NEW_ENDPOINT_KEY" + save_env_value(old_env, "old-secret-1234567890") + save_env_value(new_env, "new-secret-0987654321") + + cfg = load_config() + providers = cfg.get("providers") if isinstance(cfg.get("providers"), dict) else {} + providers[provider_id] = { + "name": "Env Preview", + "base_url": "https://env-preview.example.com/v1", + "model": "m", + "key_env": old_env, + "models": {"m": {}}, + } + cfg["providers"] = providers + save_config(cfg) + + endpoint = next( + e for e in self.client.get("/api/providers/custom-endpoints").json()["endpoints"] + if e["id"] == provider_id + ) + stale_preview = endpoint["api_key_preview"] + assert stale_preview == f"${{{old_env}}}" + + cfg = load_config() + cfg["providers"][provider_id]["key_env"] = new_env + save_config(cfg) + + response = self.client.post( + "/api/providers/custom-endpoints", + json={ + "id": provider_id, + "name": "Env Preview", + "base_url": "https://env-preview.example.com/v1", + "model": "m", + "api_key": stale_preview, + }, + ) + assert response.status_code == 400 + assert load_config()["providers"][provider_id]["key_env"] == new_env + + def test_legacy_custom_endpoint_rejects_stale_preview_after_rotation(self): + from hermes_cli.config import load_config, save_config + + provider_id = "legacy-preview" + cfg = load_config() + providers = cfg.get("providers") if isinstance(cfg.get("providers"), dict) else {} + providers[provider_id] = { + "name": "Legacy Preview", + "base_url": "https://legacy-preview.example.com/v1", + "model": "m", + "api_key": "legacy-secret-A-1234567890", + "models": {"m": {}}, + } + cfg["providers"] = providers + save_config(cfg) + + endpoint = next( + e for e in self.client.get("/api/providers/custom-endpoints").json()["endpoints"] + if e["id"] == provider_id + ) + preview = endpoint["api_key_preview"] + assert preview.startswith("«redacted") + + cfg = load_config() + cfg["providers"][provider_id]["api_key"] = "legacy-secret-B-0987654321" + save_config(cfg) + + response = self.client.post( + "/api/providers/custom-endpoints", + json={ + "id": provider_id, + "name": "Legacy Preview", + "base_url": "https://legacy-preview.example.com/v1", + "model": "m", + "api_key": preview, + }, + ) + assert response.status_code == 400 + assert load_config()["providers"][provider_id]["api_key"] == "legacy-secret-B-0987654321" + def test_messaging_rejects_preview_from_process_env(self, monkeypatch): from hermes_cli.config import load_env @@ -2183,11 +2277,15 @@ CONFIG_SCHEMA = ProviderConfigSchema( platforms = self.client.get("/api/messaging/platforms").json()["platforms"] discord = next(p for p in platforms if p["id"] == "discord") preview = next(v["redacted_value"] for v in discord["env_vars"] if v["key"] == key) - response = self.client.put(f"/api/messaging/platforms/discord", json={"env": {key: preview}}) + assert preview.startswith("«redacted") + + rotated = "discord-rotated-secret-0987654321" + monkeypatch.setenv(key, rotated) + response = self.client.put("/api/messaging/platforms/discord", json={"env": {key: preview}}) assert response.status_code == 400 assert key not in load_env() - assert os.environ[key] == real + assert os.environ[key] == rotated def test_messaging_rejects_preview_before_clearing_secret(self): from hermes_cli.config import load_env, save_env_value From af4c9286d5ae99b34246b4cc804894c651328fbc Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 13:54:42 +0530 Subject: [PATCH 04/22] fix(dashboard): reject legacy credential masks by shape, not equality is_redacted_credential_preview kept a compat branch that authorised a bare legacy mask by comparing it to the CURRENT secret. That is not provenance: GET A -> rotate to B -> submit redact_key(A) passed the check and wrote the 11-char mask over B, the same #121002 corruption class during the mixed-client window (andrexibiza's review finding on #121051). Recognise legacy masks by the producer shape of agent.redact.mask_secret (`xxxx...xxxx` or `***`) instead, drop the `current` parameter, and collapse the redundant triple check in _write_custom_endpoint to the helper plus the `${KEY_ENV}` placeholder regex. Messaging no longer needs the shared env-source helper for preflight, so the GET closure is restored. Co-authored-by: JoaoMarcos44 Co-authored-by: KoNit-K --- hermes_cli/web_routers/_common.py | 17 ++++++++++------- hermes_cli/web_routers/config_env.py | 14 ++++---------- hermes_cli/web_routers/messaging.py | 18 ++++++------------ 3 files changed, 20 insertions(+), 29 deletions(-) diff --git a/hermes_cli/web_routers/_common.py b/hermes_cli/web_routers/_common.py index 67b44b2556..e5a28a5352 100644 --- a/hermes_cli/web_routers/_common.py +++ b/hermes_cli/web_routers/_common.py @@ -7,6 +7,7 @@ from __future__ import annotations import asyncio import contextlib import logging +import re import sqlite3 import time from typing import Any, Callable, Dict, Optional @@ -131,16 +132,18 @@ def redacted_credential_preview(value: Any) -> Optional[str]: return f"«redacted:{redact_key(str(value))}»" -def is_redacted_credential_preview(submitted: Any, current: Any = None) -> bool: - """Recognize current and stale dashboard previews without guessing from key shape.""" +# Legacy bare masks (pre-sentinel pages, older Desktop builds) are recognised by the +# producer shape of ``agent.redact.mask_secret`` — never by equality to the current +# secret, which would authorise a stale preview after a rotation (#121002). +_LEGACY_MASK_RE = re.compile(r".{4}\.\.\..{4}") + + +def is_redacted_credential_preview(submitted: Any) -> bool: + """Recognize current, stale and legacy dashboard previews by shape alone.""" value = str(submitted or "") if value == "«redacted-secret»" or (value.startswith("«redacted:") and value.endswith("»")): return True - if current: - # Compatibility with a page opened before the non-reusable sentinel contract. - from hermes_cli.config import redact_key - return value == redact_key(str(current)) - return False + return value == "***" or _LEGACY_MASK_RE.fullmatch(value) is not None # Corrupt-store reporting for polled read endpoints. The dashboard polls analytics every few diff --git a/hermes_cli/web_routers/config_env.py b/hermes_cli/web_routers/config_env.py index 0487c81147..296c0d4c9a 100644 --- a/hermes_cli/web_routers/config_env.py +++ b/hermes_cli/web_routers/config_env.py @@ -299,8 +299,7 @@ async def set_env_var(body: EnvVarUpdate, profile: Optional[str] = None): from hermes_cli.credential_lifecycle import save_provider_env_credential def _save(): - current = load_env().get(body.key) - if is_redacted_credential_preview(body.value, current): + if is_redacted_credential_preview(body.value): raise ValueError(REDACTED_CREDENTIAL_WRITE_DETAIL) return save_provider_env_credential(body.key, body.value) @@ -627,14 +626,9 @@ def _write_custom_endpoint(cfg: Dict[str, Any], body: CustomEndpointUpdate) -> T env_var = custom_endpoint_key_env(endpoint_id) submitted_key = body.api_key.strip() if body.api_key is not None else None if submitted_key: - display_preview = _api_key_display(existing)[1] - existing_key_env = str(existing.get("key_env") or "").strip() - stored_secret = load_env().get(existing_key_env) if existing_key_env else None - if ( - submitted_key == display_preview - or re.fullmatch(r"\$\{[^}]+\}", submitted_key) - or is_redacted_credential_preview(submitted_key, stored_secret) - ): + # ``${KEY_ENV}`` is the GET display for key_env entries; the helper covers the + # sentinel and legacy masks. Either one is display-only, current or stale. + if re.fullmatch(r"\$\{[^}]+\}", submitted_key) or is_redacted_credential_preview(submitted_key): raise HTTPException(status_code=400, detail=REDACTED_CREDENTIAL_WRITE_DETAIL) save_env_value(env_var, submitted_key) entry["key_env"] = env_var diff --git a/hermes_cli/web_routers/messaging.py b/hermes_cli/web_routers/messaging.py index 41f146d992..3f93b66e88 100644 --- a/hermes_cli/web_routers/messaging.py +++ b/hermes_cli/web_routers/messaging.py @@ -199,11 +199,6 @@ def _platform_enablement( return enabled, configured, home_channel -def _messaging_env_value(key: str, env_on_disk: dict[str, str], *, scoped: bool) -> str: - """Resolve the credential source shared by Messaging GET and write preflight.""" - return env_on_disk.get(key) or ("" if scoped else os.getenv(key, "")) - - def _messaging_platform_payload( entry: dict[str, Any], env_on_disk: dict[str, str], runtime: dict | None, scoped: bool = False, profile_home: Optional[Path] = None, @@ -233,7 +228,9 @@ def _messaging_platform_payload( runtime_platform = {} def env_value(key: str) -> str: - return _messaging_env_value(key, env_on_disk, scoped=scoped) + # Profile-scoped: judge only the profile's own .env — the dashboard process's + # os.environ carries the ROOT install's .env and would report root credentials as the profile's. + return env_on_disk.get(key) or ("" if scoped else os.getenv(key, "")) env_vars = [ { @@ -877,12 +874,10 @@ async def update_messaging_platform(platform_id: str, body: MessagingPlatformUpd raise HTTPException(status_code=400, detail=f"{key} is not configurable for {entry['name']}") def _apply(): - with _profile_scope(target_profile) as scoped_dir: - env_on_disk = load_env() + with _profile_scope(target_profile): updates: dict[str, str] = {} - # Validate the whole request against the same credential snapshot the GET - # response uses before clearing or replacing anything. + # Validate the whole request before clearing or replacing anything. for key in body.clear_env: _check_allowed(key) for key, value in body.env.items(): @@ -890,8 +885,7 @@ async def update_messaging_platform(platform_id: str, body: MessagingPlatformUpd trimmed = value.strip() if not trimmed: continue - current = _messaging_env_value(key, env_on_disk, scoped=scoped_dir is not None) - if is_redacted_credential_preview(trimmed, current): + if is_redacted_credential_preview(trimmed): raise HTTPException(status_code=400, detail=REDACTED_CREDENTIAL_WRITE_DETAIL) _validate_messaging_env_value(platform_id, key, trimmed) updates[key] = trimmed From a7c6803410ab2c7f90dfbca7b73024abef8fd4bf Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 13:55:20 +0530 Subject: [PATCH 05/22] test(dashboard): fold preview-writeback coverage into two invariant tests Six near-duplicate tests exercised the same invariant (a display preview never writes a secret). Keep two: /api/env with the rotation + legacy-mask race, and one test for messaging preflight-before-mutation plus the custom-endpoint ${KEY_ENV} / legacy-plaintext stale previews after rotation. Drops the process-env messaging and same-snapshot custom-endpoint variants, which are subsumed by shape-based rejection. --- tests/hermes_cli/test_web_server.py | 174 +++++++--------------------- 1 file changed, 40 insertions(+), 134 deletions(-) diff --git a/tests/hermes_cli/test_web_server.py b/tests/hermes_cli/test_web_server.py index d38c9da700..878d49f523 100644 --- a/tests/hermes_cli/test_web_server.py +++ b/tests/hermes_cli/test_web_server.py @@ -2143,6 +2143,8 @@ CONFIG_SCHEMA = ProviderConfigSchema( assert "sk-in-env" not in (endpoint["api_key_preview"] or "") def test_env_rejects_its_redacted_preview(self): + """Invariant: a GET preview (sentinel or legacy bare mask) never gains write + authority, even after another actor rotates the secret behind it.""" from hermes_cli.config import load_env, save_env_value key = "OPENAI_API_KEY" @@ -2155,154 +2157,58 @@ CONFIG_SCHEMA = ProviderConfigSchema( assert response.status_code == 400 assert load_env()[key] == real - # A preview remains display-only even if another actor rotates the secret - # after the GET that minted it. rotated = "sk-rotated-secret-0987654321" save_env_value(key, rotated) - response = self.client.put("/api/env", json={"key": key, "value": preview}) - assert response.status_code == 400 - assert load_env()[key] == rotated - - def test_custom_endpoint_rejects_current_and_legacy_previews(self): - from hermes_cli.config import custom_endpoint_key_env, get_env_value - - payload = { - "id": "proxy", "name": "Proxy", "base_url": "https://llm.example.com/v1", - "model": "m", "api_key": "sk-live-secret-abcdef1234567890", - } - created = self.client.post("/api/providers/custom-endpoints", json=payload) - endpoint = next(e for e in created.json()["endpoints"] if e["id"] == "proxy") - env_key = custom_endpoint_key_env("proxy") - real = get_env_value(env_key) - - for preview in (endpoint["api_key_preview"], redact_key(real)): - response = self.client.post( - "/api/providers/custom-endpoints", json={**payload, "api_key": preview}) + for stale in (preview, redact_key(real)): + response = self.client.put("/api/env", json={"key": key, "value": stale}) assert response.status_code == 400 - assert get_env_value(env_key) == real + assert load_env()[key] == rotated - def test_custom_endpoint_rejects_stale_key_env_placeholder_after_rotation(self): - from hermes_cli.config import load_config, save_config, save_env_value - - provider_id = "env-preview" - old_env = "OLD_ENDPOINT_KEY" - new_env = "NEW_ENDPOINT_KEY" - save_env_value(old_env, "old-secret-1234567890") - save_env_value(new_env, "new-secret-0987654321") - - cfg = load_config() - providers = cfg.get("providers") if isinstance(cfg.get("providers"), dict) else {} - providers[provider_id] = { - "name": "Env Preview", - "base_url": "https://env-preview.example.com/v1", - "model": "m", - "key_env": old_env, - "models": {"m": {}}, - } - cfg["providers"] = providers - save_config(cfg) - - endpoint = next( - e for e in self.client.get("/api/providers/custom-endpoints").json()["endpoints"] - if e["id"] == provider_id - ) - stale_preview = endpoint["api_key_preview"] - assert stale_preview == f"${{{old_env}}}" - - cfg = load_config() - cfg["providers"][provider_id]["key_env"] = new_env - save_config(cfg) - - response = self.client.post( - "/api/providers/custom-endpoints", - json={ - "id": provider_id, - "name": "Env Preview", - "base_url": "https://env-preview.example.com/v1", - "model": "m", - "api_key": stale_preview, - }, - ) - assert response.status_code == 400 - assert load_config()["providers"][provider_id]["key_env"] == new_env - - def test_legacy_custom_endpoint_rejects_stale_preview_after_rotation(self): - from hermes_cli.config import load_config, save_config - - provider_id = "legacy-preview" - cfg = load_config() - providers = cfg.get("providers") if isinstance(cfg.get("providers"), dict) else {} - providers[provider_id] = { - "name": "Legacy Preview", - "base_url": "https://legacy-preview.example.com/v1", - "model": "m", - "api_key": "legacy-secret-A-1234567890", - "models": {"m": {}}, - } - cfg["providers"] = providers - save_config(cfg) - - endpoint = next( - e for e in self.client.get("/api/providers/custom-endpoints").json()["endpoints"] - if e["id"] == provider_id - ) - preview = endpoint["api_key_preview"] - assert preview.startswith("«redacted") - - cfg = load_config() - cfg["providers"][provider_id]["api_key"] = "legacy-secret-B-0987654321" - save_config(cfg) - - response = self.client.post( - "/api/providers/custom-endpoints", - json={ - "id": provider_id, - "name": "Legacy Preview", - "base_url": "https://legacy-preview.example.com/v1", - "model": "m", - "api_key": preview, - }, - ) - assert response.status_code == 400 - assert load_config()["providers"][provider_id]["api_key"] == "legacy-secret-B-0987654321" - - def test_messaging_rejects_preview_from_process_env(self, monkeypatch): - from hermes_cli.config import load_env - - key = "DISCORD_BOT_TOKEN" - real = "discord-live-secret-abcdef1234567890" - monkeypatch.setenv(key, real) - assert key not in load_env() - - platforms = self.client.get("/api/messaging/platforms").json()["platforms"] - discord = next(p for p in platforms if p["id"] == "discord") - preview = next(v["redacted_value"] for v in discord["env_vars"] if v["key"] == key) - assert preview.startswith("«redacted") - - rotated = "discord-rotated-secret-0987654321" - monkeypatch.setenv(key, rotated) - response = self.client.put("/api/messaging/platforms/discord", json={"env": {key: preview}}) - - assert response.status_code == 400 - assert key not in load_env() - assert os.environ[key] == rotated - - def test_messaging_rejects_preview_before_clearing_secret(self): - from hermes_cli.config import load_env, save_env_value + def test_messaging_and_custom_endpoint_reject_stale_previews(self): + """Invariant: preview rejection runs before any mutation (messaging clear+set), + and custom-endpoint display strings (``${KEY_ENV}`` / legacy plaintext preview) + are refused even after the entry rotated underneath them.""" + from hermes_cli.config import load_config, load_env, save_config, save_env_value key = "DISCORD_BOT_TOKEN" real = "discord-live-secret-abcdef1234567890" save_env_value(key, real) - preview = redact_key(real) - response = self.client.put( "/api/messaging/platforms/discord", - json={"clear_env": [key], "env": {key: preview}}, + json={"clear_env": [key], "env": {key: redact_key(real)}}, ) - assert response.status_code == 400 assert load_env()[key] == real + save_env_value("OLD_ENDPOINT_KEY", "old-secret-1234567890") + save_env_value("NEW_ENDPOINT_KEY", "new-secret-0987654321") + cfg = load_config() + cfg["providers"] = { + "env-preview": {"name": "Env Preview", "base_url": "https://env-preview.example.com/v1", + "model": "m", "key_env": "OLD_ENDPOINT_KEY", "models": {"m": {}}}, + "legacy-preview": {"name": "Legacy Preview", "base_url": "https://legacy-preview.example.com/v1", + "model": "m", "api_key": "legacy-secret-A-1234567890", "models": {"m": {}}}, + } + save_config(cfg) + endpoints = {e["id"]: e for e in self.client.get("/api/providers/custom-endpoints").json()["endpoints"]} + assert endpoints["env-preview"]["api_key_preview"] == "${OLD_ENDPOINT_KEY}" + assert endpoints["legacy-preview"]["api_key_preview"].startswith("«redacted") + + cfg = load_config() + cfg["providers"]["env-preview"]["key_env"] = "NEW_ENDPOINT_KEY" + cfg["providers"]["legacy-preview"]["api_key"] = "legacy-secret-B-0987654321" + save_config(cfg) + for endpoint_id, base_url in (("env-preview", "https://env-preview.example.com/v1"), + ("legacy-preview", "https://legacy-preview.example.com/v1")): + response = self.client.post("/api/providers/custom-endpoints", json={ + "id": endpoint_id, "name": "x", "base_url": base_url, "model": "m", + "api_key": endpoints[endpoint_id]["api_key_preview"], + }) + assert response.status_code == 400 + providers = load_config()["providers"] + assert providers["env-preview"]["key_env"] == "NEW_ENDPOINT_KEY" + assert providers["legacy-preview"]["api_key"] == "legacy-secret-B-0987654321" + def test_activating_an_endpoint_carries_its_credential_either_way(self): """Activate must work for both key_env and pre-#69449 plaintext entries.""" from hermes_cli.config import load_config, save_config From ea478da4013d421e83613b20f3c70256dd1bfc89 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 14:14:47 +0530 Subject: [PATCH 06/22] =?UTF-8?q?fix(dashboard):=20treat=20any=20=C2=ABred?= =?UTF-8?q?acted=E2=80=A6=20value=20as=20a=20credential=20preview?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Why: the recogniser enumerated two sentinel spellings and missed the vault marker «redacted-vault-secret» that agent.redact emits for vault values, so that display string could be written back as a secret. Reuse the same prefix test agent/redact.py uses to detect already-masked output — every «redacted… value is display-only by construction; legacy masks unchanged. --- hermes_cli/web_routers/_common.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/hermes_cli/web_routers/_common.py b/hermes_cli/web_routers/_common.py index e5a28a5352..8ccc851465 100644 --- a/hermes_cli/web_routers/_common.py +++ b/hermes_cli/web_routers/_common.py @@ -141,7 +141,10 @@ _LEGACY_MASK_RE = re.compile(r".{4}\.\.\..{4}") def is_redacted_credential_preview(submitted: Any) -> bool: """Recognize current, stale and legacy dashboard previews by shape alone.""" value = str(submitted or "") - if value == "«redacted-secret»" or (value.startswith("«redacted:") and value.endswith("»")): + # Any ``«redacted…`` value is already-masked output (the same test agent.redact uses + # to skip re-masking): our ``«redacted:…»`` sentinel, ``«redacted-secret»`` and the + # vault marker ``«redacted-vault-secret»``. Then the legacy bare mask shapes. + if value.startswith("«redacted"): return True return value == "***" or _LEGACY_MASK_RE.fullmatch(value) is not None From b9a50ddf1e18695e2b3d7c74a02bb45d8e93a95f Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 14:14:47 +0530 Subject: [PATCH 07/22] refactor(dashboard): hoist /api/env preview check out of the write thread MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Why: PUT /api/env raised the preview rejection as ValueError inside the profile-scoped thread only because _env_write_errors(http_passthrough=False) would downgrade an HTTPException to a 500. The check needs neither the profile scope nor the thread, so run it up front and raise HTTPException(400) directly like the messaging and custom-endpoint sites; the original one-line lambda is restored. Also drop the ad-hoc ${…} regex copy in favour of the existing _ENV_REF_RE.fullmatch from hermes_cli.config (identical matches). Behaviour unchanged: 400 and secret intact (gate/probes/S1_e2e.py CLEAN). --- hermes_cli/web_routers/config_env.py | 17 +++++++++-------- 1 file changed, 9 insertions(+), 8 deletions(-) diff --git a/hermes_cli/web_routers/config_env.py b/hermes_cli/web_routers/config_env.py index 296c0d4c9a..b5001978a9 100644 --- a/hermes_cli/web_routers/config_env.py +++ b/hermes_cli/web_routers/config_env.py @@ -24,7 +24,7 @@ from hermes_cli.web_server_profiles import ( _approval_mode_of, _broadcast_gateway_session_info, _is_other_profile, _parse_model_entries, ) from fastapi import HTTPException, Request -from hermes_cli.config import DEFAULT_CONFIG, OPTIONAL_ENV_VARS, read_raw_config, require_readable_config_before_write, custom_endpoint_key_env, coerce_provider_id, find_provider_entry, get_compatible_custom_providers, _deep_merge +from hermes_cli.config import DEFAULT_CONFIG, OPTIONAL_ENV_VARS, read_raw_config, require_readable_config_before_write, custom_endpoint_key_env, coerce_provider_id, find_provider_entry, get_compatible_custom_providers, _ENV_REF_RE, _deep_merge from hermes_cli.config_providers import _canonical_api_mode, _custom_provider_entry_to_provider_config from hermes_cli.web_models import ConfigUpdate, EnvVarUpdate, EnvVarDelete, EnvVarReveal, CustomEndpointUpdate from typing import Any, Dict, List, Optional, Tuple @@ -295,15 +295,16 @@ async def set_env_var(body: EnvVarUpdate, profile: Optional[str] = None): # mirror still holding the previous value of this var (model.api_key / # auxiliary.*.api_key / custom_providers[*]), so a rotation can't leave a # stale higher-precedence copy that keeps authenticating with the old key. + # Display-only previews (sentinel or legacy mask) must never gain write authority. + # Checked before the error mapper: it turns HTTPException into a 500 at this site. + if is_redacted_credential_preview(body.value): + raise HTTPException(status_code=400, detail=REDACTED_CREDENTIAL_WRITE_DETAIL) with _env_write_errors("PUT /api/env failed", http_passthrough=False): from hermes_cli.credential_lifecycle import save_provider_env_credential - def _save(): - if is_redacted_credential_preview(body.value): - raise ValueError(REDACTED_CREDENTIAL_WRITE_DETAIL) - return save_provider_env_credential(body.key, body.value) - - return await scoped_to_thread(body.profile or profile, _save) + return await scoped_to_thread( + body.profile or profile, lambda: save_provider_env_credential(body.key, body.value) + ) # Live credential probes keyed by env var: (url, auth) where auth is "bearer" @@ -628,7 +629,7 @@ def _write_custom_endpoint(cfg: Dict[str, Any], body: CustomEndpointUpdate) -> T if submitted_key: # ``${KEY_ENV}`` is the GET display for key_env entries; the helper covers the # sentinel and legacy masks. Either one is display-only, current or stale. - if re.fullmatch(r"\$\{[^}]+\}", submitted_key) or is_redacted_credential_preview(submitted_key): + if _ENV_REF_RE.fullmatch(submitted_key) or is_redacted_credential_preview(submitted_key): raise HTTPException(status_code=400, detail=REDACTED_CREDENTIAL_WRITE_DETAIL) save_env_value(env_var, submitted_key) entry["key_env"] = env_var From a7279d699ec3e6737a7755c3fc782a77c23444de Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 13:53:00 +0530 Subject: [PATCH 08/22] chore: map MarkTro1969 for salvage of #120557 --- contributors/emails/mark@svavnc.com | 2 ++ 1 file changed, 2 insertions(+) create mode 100644 contributors/emails/mark@svavnc.com diff --git a/contributors/emails/mark@svavnc.com b/contributors/emails/mark@svavnc.com new file mode 100644 index 0000000000..702bb7989d --- /dev/null +++ b/contributors/emails/mark@svavnc.com @@ -0,0 +1,2 @@ +MarkTro1969 +# PR #120557 salvage From aabb7a609d13834aeec343297365c58d43aac67a Mon Sep 17 00:00:00 2001 From: Mark DiPietro Date: Wed, 23 Sep 2026 14:24:42 -0400 Subject: [PATCH 09/22] fix(agent): make compression stall fallback deterministic (cherry picked from commit 4ac18e27445b1c667f7f9165113709e0de9fb2e3) --- agent/conversation_compression.py | 44 ++++++++ .../test_compression_attempt_lifecycle.py | 57 ++++++---- ...ompression_stall_deterministic_fallback.py | 102 ++++++++++++++++++ 3 files changed, 181 insertions(+), 22 deletions(-) diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 34c0eac45c..ef84032aad 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -1256,6 +1256,29 @@ def run_compress_context_with_progress_timeout( ) if settled: handled_exit = True + # The deadline is visible to the worker as well as the host. A cooperative summary call can + # observe it, unwind, and return the unchanged snapshot just BEFORE future.result() times out. + # Treat that settled no-op exactly like the host-side stall path; otherwise whether the + # deterministic fallback runs depends on a thread-scheduling race at the deadline. + result_messages = result[0] if isinstance(result, tuple) and result else None + if stall_fallback and fence.deadline_exceeded and result_messages is messages: + if on_timeout_cause is not None: + with _swallow('compress_context timeout-cause callback failed', exc_info=True): + on_timeout_cause(True, fence.progress_observed) + recovered = _retry_compression_on_fallback_chain( + worker=fallback_worker or worker, messages=messages, + system_prompt_fallback=system_prompt_fallback, + idle_timeout_seconds=idle, total_ceiling_seconds=ceiling, + on_commit_overrun=on_commit_overrun, on_timeout_cause=on_timeout_cause, + telemetry_agent=telemetry_agent, new_fence=new_fence, + escalate_deterministic=escalate_deterministic, + ) + if recovered is not None: + return recovered + if on_timeout is not None: + waited = time.monotonic() - wait_started + with _swallow('compress_context timeout callback failed', exc_info=True): + on_timeout(idle, waited, fence.seconds_since_progress()) return result # F6: a not-yet-started future must not linger as a stale queued job. @@ -1282,6 +1305,27 @@ def run_compress_context_with_progress_timeout( future, ceiling=ceiling, wait_started=wait_started, on_commit_overrun=on_commit_overrun ) handled_exit = True + # The cancelled worker can race the host into its commit section while unwinding a stalled + # summary. When that commit is only the unchanged snapshot, returning it here skips the + # stall-fallback ladder entirely (the over-window first-stall test then flakes). The worker is + # settled and its lease is free at this point, so retry exactly as the pre-commit cancel path + # does. A real compression result remains authoritative and returns immediately. + result_messages = result[0] if isinstance(result, tuple) and result else None + if stall_fallback and result_messages is messages: + recovered = _retry_compression_on_fallback_chain( + worker=fallback_worker or worker, messages=messages, + system_prompt_fallback=system_prompt_fallback, + idle_timeout_seconds=idle, total_ceiling_seconds=ceiling, + on_commit_overrun=on_commit_overrun, on_timeout_cause=on_timeout_cause, + telemetry_agent=telemetry_agent, new_fence=new_fence, + escalate_deterministic=escalate_deterministic, + ) + if recovered is not None: + return recovered + if on_timeout is not None: + waited = time.monotonic() - wait_started + with _swallow('compress_context timeout callback failed', exc_info=True): + on_timeout(idle, waited, fence.seconds_since_progress()) return result # Idle-timeout: cancel won pre-commit. Also free the worker's durable lease via diff --git a/tests/agent/test_compression_attempt_lifecycle.py b/tests/agent/test_compression_attempt_lifecycle.py index ff90698235..cbc042dbd9 100644 --- a/tests/agent/test_compression_attempt_lifecycle.py +++ b/tests/agent/test_compression_attempt_lifecycle.py @@ -26,6 +26,7 @@ import time from pathlib import Path from unittest.mock import patch +import agent.conversation_compression as cc from agent.conversation_compression import ( CompressionCommitFence, @@ -126,12 +127,13 @@ class TestWorkerTeardownOnCeiling: does NOT fire while the worker is alive (no overlap window).""" original = [{"role": "user", "content": "keep"}] release = threading.Event() + worker_started = threading.Event() worker_finished = threading.Event() lock_released: list[float] = [] def stuck_worker(fence: CompressionCommitFence): - # Continuous progress so only the TOTAL ceiling can expire - # (the #97488 'last progress 0.0s ago' shape). + fence.touch_progress() + worker_started.set() while not release.wait(timeout=0.02): fence.touch_progress() worker_finished.set() @@ -146,26 +148,37 @@ class TestWorkerTeardownOnCeiling: fence.register_cancelled_lock_release( lambda: lock_released.append(time.monotonic()) ) - msgs, prompt = run_compress_context_with_progress_timeout( - worker=stuck_worker, - messages=original, - system_prompt_fallback="fallback", - idle_timeout_seconds=0.1, - total_ceiling_seconds=0.3, - fence=fence, - stall_fallback=False, - ) - # Precondition: the worker is genuinely still running. - assert not worker_finished.is_set() - assert msgs is original and prompt == "fallback" - # Total-ceiling path: lease retained until the worker exits, so no - # new attempt can overlap the unchanged session. - assert not lock_released, ( - "durable lease released while the timed-out worker was still " - "alive — overlap window reopened (#97488)" - ) - release.set() - assert worker_finished.wait(timeout=2) + real_await = cc._await_worker_within_budget + + def await_after_worker_start(future, worker_fence, **kwargs): + assert worker_started.wait(timeout=1.0) + return real_await(future, worker_fence, **kwargs) + + try: + with patch.object( + cc, + "_await_worker_within_budget", + side_effect=await_after_worker_start, + ): + msgs, prompt = run_compress_context_with_progress_timeout( + worker=stuck_worker, + messages=original, + system_prompt_fallback="fallback", + idle_timeout_seconds=2.0, + total_ceiling_seconds=0.3, + fence=fence, + stall_fallback=False, + ) + assert not worker_finished.is_set() + assert msgs is original and prompt == "fallback" + assert not lock_released, ( + "durable lease released while the timed-out worker was still " + "alive — overlap window reopened (#97488)" + ) + finally: + release.set() + + assert worker_finished.wait(timeout=1.0) # Late result was fence-poisoned, never adopted. assert msgs == [{"role": "user", "content": "keep"}] diff --git a/tests/agent/test_compression_stall_deterministic_fallback.py b/tests/agent/test_compression_stall_deterministic_fallback.py index 71895bfd64..ab1d6fd889 100644 --- a/tests/agent/test_compression_stall_deterministic_fallback.py +++ b/tests/agent/test_compression_stall_deterministic_fallback.py @@ -175,6 +175,108 @@ def test_fence_level_retry_ladder_is_unchanged_without_a_prior_timeout(): release.set() +def test_deadline_observed_by_worker_still_runs_first_stall_fallback(): + """A cooperative worker can return its no-op result at the same deadline the host is waiting on. + That settled result must not win the race and bypass the deterministic over-window fallback.""" + original = [{"role": "user", "content": "keep-me"}] + recovered = [{"role": "assistant", "content": "compressed"}] + attempts = [] + timeout_causes = [] + timeout_callbacks = [] + + def primary_worker(fence: CompressionCommitFence): + attempts.append("primary") + while not fence.deadline_exceeded: + time.sleep(0.001) + return original, "unchanged" + + def fallback_worker(_fence: CompressionCommitFence): + attempts.append("fallback") + return recovered, "compressed-prompt" + + with patch("agent.auxiliary_client._get_auxiliary_task_config", return_value={"fallback_chain": []}): + msgs, prompt = run_compress_context_with_progress_timeout( + worker=primary_worker, fallback_worker=fallback_worker, messages=original, + system_prompt_fallback="degraded-prompt", idle_timeout_seconds=0.05, + total_ceiling_seconds=2.0, request_exceeds_window=True, + on_timeout_cause=lambda *cause: timeout_causes.append(cause), + on_timeout=lambda *args: timeout_callbacks.append(args), + ) + + assert msgs is recovered and prompt == "compressed-prompt" + assert attempts == ["primary", "fallback"] + assert timeout_causes == [(True, False)] + assert timeout_callbacks == [] + + +def test_admitted_unchanged_commit_at_deadline_retries_without_overlap(): + """If a worker enters its commit section at the deadline, the host must wait for that commit to + settle before minting a fresh retry fence. An unchanged commit is still a stalled attempt, not a + successful compression, so an over-window request must take the deterministic fallback rung.""" + original = [{"role": "user", "content": "keep-me"}] + recovered = [{"role": "assistant", "content": "compressed"}] + primary_fence = CompressionCommitFence() + commit_started = threading.Event() + primary_finished = threading.Event() + release_commit = threading.Event() + attempts = [] + timeout_causes = [] + timeout_callbacks = [] + retry_fences = [] + real_await = cc._await_worker_within_budget + + def primary_worker(fence: CompressionCommitFence): + attempts.append(("primary", fence)) + assert fence.begin_commit() + commit_started.set() + try: + assert release_commit.wait(timeout=2) + return original, "unchanged" + finally: + fence.finish_commit() + primary_finished.set() + + def fallback_worker(fence: CompressionCommitFence): + assert primary_finished.is_set(), "retry overlapped the admitted primary commit" + attempts.append(("fallback", fence)) + return recovered, "compressed-prompt" + + def await_at_deadline(future, fence, *, idle, ceiling, wait_started): + if fence is not primary_fence: + return real_await( + future, fence, idle=idle, ceiling=ceiling, wait_started=wait_started + ) + assert commit_started.wait(timeout=1) + # Hold the admitted commit through the deadline, then let it finish while + # _await_in_flight_commit owns the host-side wait. + time.sleep(ceiling + 0.01) + threading.Timer(0.02, release_commit.set).start() + return False, None + + def new_fence(): + fence = CompressionCommitFence() + retry_fences.append(fence) + return fence + + with patch.object(cc, "_await_worker_within_budget", side_effect=await_at_deadline), patch( + "agent.auxiliary_client._get_auxiliary_task_config", return_value={"fallback_chain": []} + ): + msgs, prompt = run_compress_context_with_progress_timeout( + worker=primary_worker, fallback_worker=fallback_worker, messages=original, + system_prompt_fallback="degraded-prompt", idle_timeout_seconds=0.05, + total_ceiling_seconds=0.05, request_exceeds_window=True, fence=primary_fence, + new_fence=new_fence, on_timeout_cause=lambda *cause: timeout_causes.append(cause), + on_timeout=lambda *args: timeout_callbacks.append(args), + ) + + assert msgs is recovered and prompt == "compressed-prompt" + assert primary_finished.is_set() + assert retry_fences and retry_fences[0] is not primary_fence + assert attempts == [("primary", primary_fence), ("fallback", retry_fences[0])] + assert timeout_causes == [(True, False)] + assert timeout_callbacks == [] + + def test_over_window_request_commits_the_deterministic_fallback_on_the_first_stall(tmp_path, fast_timeouts): """A request above the model's context window cannot be sent unchanged, so waiting for a SECOND stall (which on the messaging gateway never comes — the first one auto-reset the session) is a dead end: From f0672a0c436b2881076d6abfcec13e7549da34ed Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 13:55:14 +0530 Subject: [PATCH 10/22] refactor(compression): single stall-fallback ladder for every stall exit The settled-at-deadline no-op path, the post-cancel unchanged-commit path and the idle-timeout path each carried their own copy of the same retry-chain -> on_timeout -> degraded-prompt ladder (three copies after the deterministic-fallback fix). Fold them into one closure, _recover_from_stall, plus _is_unchanged_snapshot for the identity check. Behavioural alignment for the two newer paths: when the retry chain yields None they now return (messages, fallback prompt) like the idle path did, instead of the worker's own tuple with its unset prompt, and they log the same "no progress" warning when no on_timeout callback is installed. Co-authored-by: Mark DiPietro --- agent/conversation_compression.py | 89 ++++++++++++------------------- 1 file changed, 35 insertions(+), 54 deletions(-) diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index ef84032aad..5430fdf551 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -1250,6 +1250,36 @@ def run_compress_context_with_progress_timeout( # EVERY host unwind must revoke commit admission or a detached worker could # later mutate durable state; handled_exit marks paths that settle it themselves handled_exit = False + + def _is_unchanged_snapshot(result: Any) -> bool: + # A worker that observed the deadline/cancel returns the SAME messages object it was handed. + return isinstance(result, tuple) and bool(result) and result[0] is messages + + def _recover_from_stall(since_progress: float) -> tuple[list[dict[str, Any]], str]: + """One stall-fallback ladder for every host-side stall exit (idle timeout, settled-at-deadline + no-op, post-cancel unchanged commit): retry chain, then on_timeout, then the degraded prompt.""" + # Lease is free, so run the fallback BEFORE on_timeout: that callback records + # the summary-failure cooldown, which would no-op the retry's summary call. + if stall_fallback: + recovered = _retry_compression_on_fallback_chain( + worker=fallback_worker or worker, messages=messages, system_prompt_fallback=system_prompt_fallback, + idle_timeout_seconds=idle, total_ceiling_seconds=ceiling, on_commit_overrun=on_commit_overrun, + on_timeout_cause=on_timeout_cause, telemetry_agent=telemetry_agent, new_fence=new_fence, + escalate_deterministic=escalate_deterministic, + ) + if recovered is not None: + return recovered + waited = time.monotonic() - wait_started + if on_timeout is not None: + with _swallow('compress_context timeout callback failed', exc_info=True): + on_timeout(idle, waited, since_progress) + else: + logger.warning( + "Context compression made no progress for %.1fs (total wait %.1fs, ceiling %.1fs); continuing without " + "compression", since_progress, waited, ceiling, + ) + return messages, _resolve_fallback_prompt() + try: settled, result = _await_worker_within_budget( future, fence, idle=idle, ceiling=ceiling, wait_started=wait_started @@ -1260,25 +1290,11 @@ def run_compress_context_with_progress_timeout( # observe it, unwind, and return the unchanged snapshot just BEFORE future.result() times out. # Treat that settled no-op exactly like the host-side stall path; otherwise whether the # deterministic fallback runs depends on a thread-scheduling race at the deadline. - result_messages = result[0] if isinstance(result, tuple) and result else None - if stall_fallback and fence.deadline_exceeded and result_messages is messages: + if stall_fallback and fence.deadline_exceeded and _is_unchanged_snapshot(result): if on_timeout_cause is not None: with _swallow('compress_context timeout-cause callback failed', exc_info=True): on_timeout_cause(True, fence.progress_observed) - recovered = _retry_compression_on_fallback_chain( - worker=fallback_worker or worker, messages=messages, - system_prompt_fallback=system_prompt_fallback, - idle_timeout_seconds=idle, total_ceiling_seconds=ceiling, - on_commit_overrun=on_commit_overrun, on_timeout_cause=on_timeout_cause, - telemetry_agent=telemetry_agent, new_fence=new_fence, - escalate_deterministic=escalate_deterministic, - ) - if recovered is not None: - return recovered - if on_timeout is not None: - waited = time.monotonic() - wait_started - with _swallow('compress_context timeout callback failed', exc_info=True): - on_timeout(idle, waited, fence.seconds_since_progress()) + return _recover_from_stall(fence.seconds_since_progress()) return result # F6: a not-yet-started future must not linger as a stale queued job. @@ -1310,55 +1326,20 @@ def run_compress_context_with_progress_timeout( # stall-fallback ladder entirely (the over-window first-stall test then flakes). The worker is # settled and its lease is free at this point, so retry exactly as the pre-commit cancel path # does. A real compression result remains authoritative and returns immediately. - result_messages = result[0] if isinstance(result, tuple) and result else None - if stall_fallback and result_messages is messages: - recovered = _retry_compression_on_fallback_chain( - worker=fallback_worker or worker, messages=messages, - system_prompt_fallback=system_prompt_fallback, - idle_timeout_seconds=idle, total_ceiling_seconds=ceiling, - on_commit_overrun=on_commit_overrun, on_timeout_cause=on_timeout_cause, - telemetry_agent=telemetry_agent, new_fence=new_fence, - escalate_deterministic=escalate_deterministic, - ) - if recovered is not None: - return recovered - if on_timeout is not None: - waited = time.monotonic() - wait_started - with _swallow('compress_context timeout callback failed', exc_info=True): - on_timeout(idle, waited, fence.seconds_since_progress()) + if stall_fallback and _is_unchanged_snapshot(result): + return _recover_from_stall(fence.seconds_since_progress()) return result # Idle-timeout: cancel won pre-commit. Also free the worker's durable lease via # the holder-qualified hook so a NEW compressor can acquire at once (no ABA). handled_exit = True _release_cancelled_worker(future, fence, total_exhausted=total_exhausted, ceiling=ceiling) - waited = time.monotonic() - wait_started # #76354 S3 analogue for this wait: charge the idle budget from the LAST PROGRESS event, not from # the start of this wait slice. Waiting a full ``idle`` after progress that landed early in the # previous slice would allow silence to approach 2x the budget. - since_progress = fence.seconds_since_progress() - # Lease is free, so run the fallback BEFORE on_timeout: that callback records - # the summary-failure cooldown, which would no-op the retry's summary call. - if stall_fallback: - recovered = _retry_compression_on_fallback_chain( - worker=fallback_worker or worker, messages=messages, system_prompt_fallback=system_prompt_fallback, - idle_timeout_seconds=idle, total_ceiling_seconds=ceiling, on_commit_overrun=on_commit_overrun, - on_timeout_cause=on_timeout_cause, telemetry_agent=telemetry_agent, new_fence=new_fence, - escalate_deterministic=escalate_deterministic, - ) - if recovered is not None: - return recovered - if on_timeout is not None: - with _swallow('compress_context timeout callback failed', exc_info=True): - on_timeout(idle, waited, since_progress) - else: - logger.warning( - "Context compression made no progress for %.1fs (total wait %.1fs, ceiling %.1fs); continuing without " - "compression", since_progress, waited, ceiling, - ) # Leave the future on the shared pool: fence cancel won, so a late # commit cannot land (same detachment model as gateway hygiene). - return messages, _resolve_fallback_prompt() + return _recover_from_stall(fence.seconds_since_progress()) finally: if not handled_exit: # Any unwind while waiting: revoke commit admission and release the worker's From 3bd93daedc9ffb04ad811bdf6af0add8f1fc323d Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 14:58:07 +0530 Subject: [PATCH 11/22] refactor(compression): sample the stall wait once, before the retry chain _recover_from_stall now samples both the total wait and seconds-since-progress itself, before the fallback retry, so the "reached its total ceiling after Ns" warning reports the stall rather than stall plus retry time (as before the ladder refactor) and the three callers stop passing the same value. The snapshot check trusts the worker's Tuple[list, str] contract like the rest of the module, and the teardown test goes back to a 0.3s idle budget: the 2.0s idle made the 0.3s ceiling dead and stretched the join grace to 2s. --- agent/conversation_compression.py | 19 ++++++++++--------- .../test_compression_attempt_lifecycle.py | 4 +++- 2 files changed, 13 insertions(+), 10 deletions(-) diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 5430fdf551..45511e4914 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -1253,11 +1253,16 @@ def run_compress_context_with_progress_timeout( def _is_unchanged_snapshot(result: Any) -> bool: # A worker that observed the deadline/cancel returns the SAME messages object it was handed. - return isinstance(result, tuple) and bool(result) and result[0] is messages + return result[0] is messages - def _recover_from_stall(since_progress: float) -> tuple[list[dict[str, Any]], str]: + def _recover_from_stall() -> tuple[list[dict[str, Any]], str]: """One stall-fallback ladder for every host-side stall exit (idle timeout, settled-at-deadline no-op, post-cancel unchanged commit): retry chain, then on_timeout, then the degraded prompt.""" + # Sample before the retry chain so the reported wait is the stall itself, not stall + retry. + # #76354 S3 analogue: silence is charged from the LAST PROGRESS event, not from the start of + # this wait slice, or progress early in a previous slice would let silence approach 2x idle. + waited = time.monotonic() - wait_started + since_progress = fence.seconds_since_progress() # Lease is free, so run the fallback BEFORE on_timeout: that callback records # the summary-failure cooldown, which would no-op the retry's summary call. if stall_fallback: @@ -1269,7 +1274,6 @@ def run_compress_context_with_progress_timeout( ) if recovered is not None: return recovered - waited = time.monotonic() - wait_started if on_timeout is not None: with _swallow('compress_context timeout callback failed', exc_info=True): on_timeout(idle, waited, since_progress) @@ -1294,7 +1298,7 @@ def run_compress_context_with_progress_timeout( if on_timeout_cause is not None: with _swallow('compress_context timeout-cause callback failed', exc_info=True): on_timeout_cause(True, fence.progress_observed) - return _recover_from_stall(fence.seconds_since_progress()) + return _recover_from_stall() return result # F6: a not-yet-started future must not linger as a stale queued job. @@ -1327,19 +1331,16 @@ def run_compress_context_with_progress_timeout( # settled and its lease is free at this point, so retry exactly as the pre-commit cancel path # does. A real compression result remains authoritative and returns immediately. if stall_fallback and _is_unchanged_snapshot(result): - return _recover_from_stall(fence.seconds_since_progress()) + return _recover_from_stall() return result # Idle-timeout: cancel won pre-commit. Also free the worker's durable lease via # the holder-qualified hook so a NEW compressor can acquire at once (no ABA). handled_exit = True _release_cancelled_worker(future, fence, total_exhausted=total_exhausted, ceiling=ceiling) - # #76354 S3 analogue for this wait: charge the idle budget from the LAST PROGRESS event, not from - # the start of this wait slice. Waiting a full ``idle`` after progress that landed early in the - # previous slice would allow silence to approach 2x the budget. # Leave the future on the shared pool: fence cancel won, so a late # commit cannot land (same detachment model as gateway hygiene). - return _recover_from_stall(fence.seconds_since_progress()) + return _recover_from_stall() finally: if not handled_exit: # Any unwind while waiting: revoke commit admission and release the worker's diff --git a/tests/agent/test_compression_attempt_lifecycle.py b/tests/agent/test_compression_attempt_lifecycle.py index cbc042dbd9..3351468a41 100644 --- a/tests/agent/test_compression_attempt_lifecycle.py +++ b/tests/agent/test_compression_attempt_lifecycle.py @@ -132,6 +132,8 @@ class TestWorkerTeardownOnCeiling: lock_released: list[float] = [] def stuck_worker(fence: CompressionCommitFence): + # Continuous progress so only the TOTAL ceiling can expire + # (the #97488 'last progress 0.0s ago' shape). fence.touch_progress() worker_started.set() while not release.wait(timeout=0.02): @@ -164,7 +166,7 @@ class TestWorkerTeardownOnCeiling: worker=stuck_worker, messages=original, system_prompt_fallback="fallback", - idle_timeout_seconds=2.0, + idle_timeout_seconds=0.3, total_ceiling_seconds=0.3, fence=fence, stall_fallback=False, From 3028995817041559572eba1a2f72878481fc345c Mon Sep 17 00:00:00 2001 From: John Paul Soliva Date: Wed, 23 Sep 2026 20:38:58 +0900 Subject: [PATCH 12/22] fix(compression): an in-place compaction never archives turns the compacting surface did not hold The in-place commit archives every active row up to the lease watermark, the newest row in state.db. But a surface compacts the history it holds: the Desktop/TUI session.compress RPC compacts session["history"], the CLI's /compress its conversation_history. When another surface appended turns to the same session since (a Desktop session continued from Telegram after /handoff, #42962, or a CLI resumed elsewhere), those rows sat under the watermark, took the positional rewind slots and became active=0, compacted=0: gone from every surface's model and display history, REST, the gateway's next replay and session_search. The summarizer never saw them, and nothing re-adopts them. `/compress here N` (#119962) loses them the same way. The watermark is now capped at the newest durable row the compressor was handed (its messages plus the kept here-N tail). Rows above it take the existing concurrent-append path, cloned after the compacted set. That is the premise _adopt_out_of_band_turns already rests on, so the next prompt adopts them as before. The cap applies only while the held history is a live prefix of the session: its newest durable row must carry its _row_id and still be active. Otherwise (a history another surface already compacted, or rows held without ids, as the gateway's replay dicts are) the commit falls back to the lease watermark, unchanged. (cherry picked from commit 930866a2bbc2ec460898a0174c3600d97d8ed088) --- agent/conversation_compression.py | 33 +++++++++- .../test_conversation_compression_manual.py | 61 +++++++++++++++++-- 2 files changed, 88 insertions(+), 6 deletions(-) diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 45511e4914..1e85ee1cde 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -3607,6 +3607,35 @@ class _CommitOutcome: made_progress: bool = False +def _held_watermark(agent: Any, watermark: Optional[int], messages: list, verbatim_tail: Optional[list]) -> Optional[int]: + """The in-place archive watermark, capped at the newest durable row the compressor was handed. + + The lease watermark is the newest row in state.db, but a surface compacts the history it holds, and that + can be older: a Desktop/TUI or CLI /compress, or a long-lived CLI, does not hold turns another surface + appended to the same session since. Archived under the watermark, those rows would leave every surface's + history and search, and the summary never saw them. Above the cap they take the concurrent-append path + instead (cloned after the compacted set). + + Only while the held history is a live prefix of the session: its newest durable row names its ``_row_id`` + (one held without it could sit above the cap and be cloned beside its own carried copy), and that row is + still active (after another surface compacted, the held rows are archived and every live row would be + cloned beside the new summary). + """ + if watermark is None: + return None + from agent.context_compressor import _DB_PERSISTED_MARKER + rows = [*((m, False) for m in messages), *((m, True) for m in verbatim_tail or ())] + held = [m.get("_row_id") for m, _ in rows if isinstance(m, dict)] + held = [rid for rid in held if isinstance(rid, int) and not isinstance(rid, bool) and rid > 0] + newest = next((m for m, kept in reversed(rows) if isinstance(m, dict) + and (kept or "_row_id" in m or m.get(_DB_PERSISTED_MARKER))), None) + if newest is None or newest.get("_row_id") not in held or max(held) >= watermark: + return watermark + if agent._session_db.get_message_role(agent.session_id, max(held)) is None: + return watermark + return max(held) + + def _commit_compaction( agent: Any, messages: list, compressed: list, *, in_place: bool, lease: _CompressionLease, new_system_prompt: str, system_message: str, compressed_user_turn_outcome: str, @@ -3673,8 +3702,8 @@ def _commit_compaction( tail_count += len(verbatim_tail) agent._session_db.archive_and_compact( agent.session_id, persisted, model_config_patch={PROACTIVE_PRUNE_REARM_MODEL_CONFIG_KEY: None}, - watermark=lease.watermark, lock_holder=lease.holder, tail_count=tail_count, - carried_messages=carried_messages, + watermark=_held_watermark(agent, lease.watermark, messages, verbatim_tail), + lock_holder=lease.holder, tail_count=tail_count, carried_messages=carried_messages, ) compressed = persisted split_status = "in_place_committed" diff --git a/tests/agent/test_conversation_compression_manual.py b/tests/agent/test_conversation_compression_manual.py index 1808302305..c55e6c65ab 100644 --- a/tests/agent/test_conversation_compression_manual.py +++ b/tests/agent/test_conversation_compression_manual.py @@ -130,9 +130,10 @@ def _exchanges(n, *, unanswered=()): return history -def _stored_agent(db, history): +def _stored_agent(db, history, *, create=True): """A real AIAgent (default in-place mode) over ``history`` as stored rows, loaded back like the gateway does.""" - db.create_session("sid", "telegram", model="test/model") + if create: + db.create_session("sid", "telegram", model="test/model") for message in history: db.append_message("sid", message["role"], message["content"]) with patch.dict(os.environ, {"OPENROUTER_API_KEY": "test-key"}): @@ -144,12 +145,15 @@ def _stored_agent(db, history): def _compress_here(agent, history, keep): + return _compress(agent, history, f"here {keep}") + + +def _compress(agent, history, raw): response = MagicMock() response.choices = [MagicMock()] response.choices[0].message.content = "## Goal\nNumbered fruit questions.\n## Progress\nEarly ones answered." with patch("agent.context_compressor.call_llm", lambda **_kw: response): - return compress_now(agent, history, parse_compress_args(f"here {keep}"), system_message="", - skip_without_window=True) + return compress_now(agent, history, parse_compress_args(raw), system_message="", skip_without_window=True) def _flags(db, content): @@ -208,3 +212,52 @@ def test_in_place_here_n_folds_the_seam_once(session_db): assert _live(durable) == _live(result.after_messages) assert f"{history[-5]['content']}\n\n{history[-4]['content']}" in [m["content"] for m in durable] assert _live(durable[-3:]) == _live(history[-3:]) + + +FOREIGN_TURN = [("user", "[from Telegram] the vault code is 7741"), ("assistant", "Noted: the vault code is 7741.")] + + +@pytest.mark.parametrize("raw", ["", "here 2"]) +def test_in_place_compress_keeps_turns_the_caller_never_held(session_db, raw): + """A surface compacts the history it holds (Desktop/TUI and CLI /compress), and another surface may have + appended turns to the same session since. The archive runs under state.db's newest row, so unless it stops at + the caller's newest row it takes those turns away unseen: the summarizer never read them, and they leave every + surface's history and search.""" + agent, _ = _stored_agent(session_db, _exchanges(10)) + held = session_db.get_resume_conversations("sid")[0] # what a resume restores: row ids included + for role, content in FOREIGN_TURN: + session_db.append_message("sid", role, content) + + assert _compress(agent, held, raw).status == "compressed" + + durable = session_db.get_messages_as_conversation("sid") + assert _live(durable[-2:]) == FOREIGN_TURN + assert len({m["content"] for m in durable}) == len(durable) # kept (`here 2`) and carried rows once each + model_history, display_history = session_db.get_resume_conversations("sid") + for _role, content in FOREIGN_TURN: + assert [m["content"] for m in model_history].count(content) == 1 + assert [m["content"] for m in display_history].count(content) == 1 + assert session_db.search_messages("vault 7741") + + +@pytest.mark.parametrize("stale", ["compacted elsewhere", "newest rows without ids"]) +def test_in_place_compress_never_leaves_two_live_copies_of_a_row(session_db, stale): + """Stopping the archive at the caller's newest row is only exact while the held history is a live prefix of the + session. After another surface compacted it, every live row is newer than the held (now archived) ones; and a + durable row held without its row id may sit above the stop. Either way the rows above it would be cloned beside + their own copies in the new transcript.""" + agent, _ = _stored_agent(session_db, _exchanges(10)) + held = session_db.get_resume_conversations("sid")[0] + if stale == "compacted elsewhere": + other, _ = _stored_agent(session_db, [], create=False) + assert _compress(other, session_db.get_messages_as_conversation("sid"), "").status == "compressed" + else: + for message in held[-2:]: + message.pop("_row_id") + for role, content in FOREIGN_TURN: + session_db.append_message("sid", role, content) + + assert _compress(agent, held, "").status == "compressed" + + contents = [m["content"] for m in session_db.get_messages_as_conversation("sid")] + assert len(contents) == len(set(contents)) From f53c0f6b20a97d401b4e75325b8e981cfcad68e8 Mon Sep 17 00:00:00 2001 From: teknium1 <127238744+teknium1@users.noreply.github.com> Date: Wed, 23 Sep 2026 05:36:44 -0700 Subject: [PATCH 13/22] fix(compression): a trailing held row without any stamp keeps the lease watermark _held_watermark skipped trailing dicts that carried neither `_row_id` nor the persisted marker when picking the "newest held row", assuming such a row is not durable. The TUI model-switch marker is: server.py appends the bare dict to session["history"] and writes it with a plain append_message, stamping nothing. With that row last in history a plain /compress capped the archive at the previous stamped row, so the marker's durable row sat above the cap, was cloned as a "concurrent append" AND inserted from the compacted set: two live copies (probe on the PR head: rows 32 and 33 both the marker; origin/main keeps one). Any trailing dict of unknown provenance now disables the cap (lease watermark, today's behaviour) instead of being looked past. The regression case rides the existing never-leaves-two-live-copies test as a third parameter; red on the previous predicate. Review-fix on #120156. (cherry picked from commit b114f564f0eddb3ccfe16123e364097747562928) --- agent/conversation_compression.py | 17 ++++++++--------- .../test_conversation_compression_manual.py | 14 ++++++++++---- 2 files changed, 18 insertions(+), 13 deletions(-) diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 1e85ee1cde..334c31b20c 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -3616,19 +3616,18 @@ def _held_watermark(agent: Any, watermark: Optional[int], messages: list, verbat history and search, and the summary never saw them. Above the cap they take the concurrent-append path instead (cloned after the compacted set). - Only while the held history is a live prefix of the session: its newest durable row names its ``_row_id`` - (one held without it could sit above the cap and be cloned beside its own carried copy), and that row is - still active (after another surface compacted, the held rows are archived and every live row would be - cloned beside the new summary). + Only while the held history is a live prefix of the session: its LAST row names its ``_row_id`` (a + trailing row of unknown provenance may be durable under the lease watermark without any stamp, e.g. the + TUI model-switch marker appended to history and written with a bare ``append_message``; capped below it, + the clone would land beside its own carried copy), and that row is still active (after another surface + compacted, the held rows are archived and every live row would be cloned beside the new summary). """ if watermark is None: return None - from agent.context_compressor import _DB_PERSISTED_MARKER - rows = [*((m, False) for m in messages), *((m, True) for m in verbatim_tail or ())] - held = [m.get("_row_id") for m, _ in rows if isinstance(m, dict)] + rows = [*messages, *(verbatim_tail or ())] + held = [m.get("_row_id") for m in rows if isinstance(m, dict)] held = [rid for rid in held if isinstance(rid, int) and not isinstance(rid, bool) and rid > 0] - newest = next((m for m, kept in reversed(rows) if isinstance(m, dict) - and (kept or "_row_id" in m or m.get(_DB_PERSISTED_MARKER))), None) + newest = next((m for m in reversed(rows) if isinstance(m, dict)), None) if newest is None or newest.get("_row_id") not in held or max(held) >= watermark: return watermark if agent._session_db.get_message_role(agent.session_id, max(held)) is None: diff --git a/tests/agent/test_conversation_compression_manual.py b/tests/agent/test_conversation_compression_manual.py index c55e6c65ab..95e694b915 100644 --- a/tests/agent/test_conversation_compression_manual.py +++ b/tests/agent/test_conversation_compression_manual.py @@ -240,17 +240,23 @@ def test_in_place_compress_keeps_turns_the_caller_never_held(session_db, raw): assert session_db.search_messages("vault 7741") -@pytest.mark.parametrize("stale", ["compacted elsewhere", "newest rows without ids"]) +@pytest.mark.parametrize("stale", ["compacted elsewhere", "newest rows without ids", "trailing row without any stamp"]) def test_in_place_compress_never_leaves_two_live_copies_of_a_row(session_db, stale): """Stopping the archive at the caller's newest row is only exact while the held history is a live prefix of the - session. After another surface compacted it, every live row is newer than the held (now archived) ones; and a - durable row held without its row id may sit above the stop. Either way the rows above it would be cloned beside - their own copies in the new transcript.""" + session. After another surface compacted it, every live row is newer than the held (now archived) ones; a + durable row held without its row id may sit above the stop; and a trailing row held with neither a row id nor + the persisted marker can still be durable under the lease watermark (the TUI model-switch marker is appended + to history and written with a bare ``append_message``). Either way the rows above the stop would be cloned + beside their own copies in the new transcript.""" agent, _ = _stored_agent(session_db, _exchanges(10)) held = session_db.get_resume_conversations("sid")[0] if stale == "compacted elsewhere": other, _ = _stored_agent(session_db, [], create=False) assert _compress(other, session_db.get_messages_as_conversation("sid"), "").status == "compressed" + elif stale == "trailing row without any stamp": + marker = "[Model switched to test/other.]" + held.append({"role": "user", "content": marker, "display_kind": "model_switch"}) + session_db.append_message("sid", "user", marker, display_kind="model_switch") else: for message in held[-2:]: message.pop("_row_id") From 1a632c3ab3b2108135fc1d15c8e9b9031cd4289c Mon Sep 17 00:00:00 2001 From: John Paul Soliva Date: Thu, 24 Sep 2026 04:26:19 +0900 Subject: [PATCH 14/22] fix(compression): a held row a merge rewrote does not bound the archive Resume repair merges consecutive user (and assistant) rows into the first one's dict: that dict keeps its _row_id, drops the persisted marker, and the later row's id leaves the held history. _held_watermark capped the archive at max(held), so the later row sat above the cap, was cloned as a concurrent append, and stayed live beside the merged row that already carries its content. A context engine that rewrites the held list in place reaches the same state. A dict loaded from the DB is born carrying both the id and the marker, so a newest held row with an id but no marker is exactly the rewritten case: keep the lease watermark there, as main does. A here-N tail is exempt; it is copied verbatim without the marker, its ids are exact, and being newest they lift the cap above anything a merge earlier in the history absorbed. The new test reproduces the merge through get_resume_conversations, the real resume path, rather than a hand-built dict. The existing exact-duplicate check cannot see this case (the merged and cloned rows are different strings), so it asserts the absorbed prompt appears in one live row. Red on the previous head, green here. (cherry picked from commit c41b9df15d72765d7d4fbb690077de463d37beed) --- agent/conversation_compression.py | 11 +++++++++++ .../test_conversation_compression_manual.py | 17 +++++++++++++++++ 2 files changed, 28 insertions(+) diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 334c31b20c..02e60f7fdf 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -3624,12 +3624,23 @@ def _held_watermark(agent: Any, watermark: Optional[int], messages: list, verbat """ if watermark is None: return None + from agent.context_compressor import _DB_PERSISTED_MARKER rows = [*messages, *(verbatim_tail or ())] held = [m.get("_row_id") for m in rows if isinstance(m, dict)] held = [rid for rid in held if isinstance(rid, int) and not isinstance(rid, bool) and rid > 0] newest = next((m for m in reversed(rows) if isinstance(m, dict)), None) if newest is None or newest.get("_row_id") not in held or max(held) >= watermark: return watermark + # ...and, when it is a summarized row, that it still matches its durable row. A dict loaded from + # the DB is born carrying both the id and the persist marker; a pass that rewrote its content drops + # the marker and keeps the id (the user/assistant merges in `repair_message_sequence`, or a context + # engine that rewrites in place). Such a row absorbed later durable rows whose ids are gone from + # `held`, so `max(held)` no longer names what the summary covered: capped there, those rows would be + # cloned live beside a summary that already contains them. A `here N` tail is exempt: it is copied + # verbatim (without the marker), so its ids are exact and, being newest, they lift the cap above + # anything a merged row earlier in `messages` absorbed. + if not any(isinstance(m, dict) for m in (verbatim_tail or ())) and not newest.get(_DB_PERSISTED_MARKER): + return watermark if agent._session_db.get_message_role(agent.session_id, max(held)) is None: return watermark return max(held) diff --git a/tests/agent/test_conversation_compression_manual.py b/tests/agent/test_conversation_compression_manual.py index 95e694b915..3e81880cf4 100644 --- a/tests/agent/test_conversation_compression_manual.py +++ b/tests/agent/test_conversation_compression_manual.py @@ -267,3 +267,20 @@ def test_in_place_compress_never_leaves_two_live_copies_of_a_row(session_db, sta contents = [m["content"] for m in session_db.get_messages_as_conversation("sid")] assert len(contents) == len(set(contents)) + + +def test_in_place_compress_never_clones_a_row_a_merge_already_carried(session_db): + """Resume repair merges consecutive user rows into the first one's dict: that dict keeps its row id, drops the + persisted marker, and the later row's id leaves the held history. Stopping the archive at the newest held id + would then leave the later row above the stop, cloned as a concurrent append beside the merged row that + already carries its content. Two strings, so the exact-duplicate check above cannot see it.""" + agent, _ = _stored_agent(session_db, _exchanges(10)) + session_db.append_message("sid", "user", "first half of a split prompt") + session_db.append_message("sid", "user", "second half 9931") + held = session_db.get_resume_conversations("sid")[0] # the resume path merges the two user rows + assert "second half 9931" in held[-1]["content"] and held[-1]["content"] != "second half 9931" + + assert _compress(agent, held, "").status == "compressed" + + contents = [m["content"] for m in session_db.get_messages_as_conversation("sid")] + assert sum("second half 9931" in c for c in contents) == 1 From 4d5899868a5e5b0af54ee7d132ba505c38d25ded Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 14:43:30 +0530 Subject: [PATCH 15/22] fix(compression): a here-N tail copy vouches for its row id only when its source is unchanged `/compress here N` left two live copies of a resume-merged user row when the merged row was the newest tail row. Resume repair merges consecutive user rows into the first dict (keeps its _row_id, drops the persist marker; the second row's id leaves the held set). _held_watermark exempted a verbatim tail from the marker check because its copies are marker-swept by construction, assuming their ids were exact and newest. That was false for the merged case: the cap landed on the merged row's id and the absorbed durable row above it was cloned as a concurrent append beside the row that already carries its text (main archives it). Put the exactness decision where the marker is still visible: compress_now keeps a tail copy's _row_id only when the source dict still carries the marker. _held_watermark then applies one rule to every row (an id is exact when its dict is unchanged since load: marker on a held row, a tail copy's id trusted as given) and the newest row must have one, with no tail exemption and max(held) computed once. `here 2` joins the existing merge test; it is red without this change. --- agent/conversation_compression.py | 51 ++++++++++--------- agent/conversation_compression_manual.py | 11 +++- .../test_conversation_compression_manual.py | 8 +-- 3 files changed, 41 insertions(+), 29 deletions(-) diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 02e60f7fdf..209907df7b 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -3611,39 +3611,42 @@ def _held_watermark(agent: Any, watermark: Optional[int], messages: list, verbat """The in-place archive watermark, capped at the newest durable row the compressor was handed. The lease watermark is the newest row in state.db, but a surface compacts the history it holds, and that - can be older: a Desktop/TUI or CLI /compress, or a long-lived CLI, does not hold turns another surface - appended to the same session since. Archived under the watermark, those rows would leave every surface's + can be older: a Desktop/TUI or CLI ``/compress`` does not hold turns another surface appended to the same + session since it loaded. Archived under the watermark, those rows would leave every surface's history and search, and the summary never saw them. Above the cap they take the concurrent-append path instead (cloned after the compacted set). - Only while the held history is a live prefix of the session: its LAST row names its ``_row_id`` (a - trailing row of unknown provenance may be durable under the lease watermark without any stamp, e.g. the - TUI model-switch marker appended to history and written with a bare ``append_message``; capped below it, - the clone would land beside its own carried copy), and that row is still active (after another surface - compacted, the held rows are archived and every live row would be cloned beside the new summary). + Only while the held history is a live prefix of the session: its LAST row names an exact ``_row_id``, and + that row is still active (after another surface compacted, the held rows are archived and every live row + would be cloned beside the new summary). An id is exact while its dict is unchanged since it was loaded: a + dict loaded from the DB is born carrying both the id and the persist marker, and a pass that rewrote its + content drops the marker and keeps the id (the user/assistant merges in ``repair_message_sequence``, or a + context engine that rewrites in place). Such a row absorbed later durable rows whose ids are gone from the + held set, so its id no longer names what the summary covered: capped there, those rows would be cloned + live beside a summary that already contains them. A trailing row of unknown provenance may be durable + under the lease watermark without any stamp (the TUI model-switch marker, written with a bare + ``append_message``); capped below it, the clone would land beside its own carried copy. A ``here N`` tail + is marker-swept copies, so ``compress_now`` keeps a copy's id only when its source still carried the + marker; their ids are trusted as given. """ if watermark is None: return None from agent.context_compressor import _DB_PERSISTED_MARKER - rows = [*messages, *(verbatim_tail or ())] - held = [m.get("_row_id") for m in rows if isinstance(m, dict)] - held = [rid for rid in held if isinstance(rid, int) and not isinstance(rid, bool) and rid > 0] - newest = next((m for m in reversed(rows) if isinstance(m, dict)), None) - if newest is None or newest.get("_row_id") not in held or max(held) >= watermark: + + def _exact_id(m: dict, copied: bool) -> Optional[int]: + rid = m.get("_row_id") + if not isinstance(rid, int) or isinstance(rid, bool) or rid <= 0: + return None + return rid if copied or m.get(_DB_PERSISTED_MARKER) else None + + ids = [_exact_id(m, False) for m in messages if isinstance(m, dict)] + ids += [_exact_id(m, True) for m in (verbatim_tail or ()) if isinstance(m, dict)] + held = [rid for rid in ids if rid is not None] + if not ids or ids[-1] is None or (newest_held := max(held)) >= watermark: return watermark - # ...and, when it is a summarized row, that it still matches its durable row. A dict loaded from - # the DB is born carrying both the id and the persist marker; a pass that rewrote its content drops - # the marker and keeps the id (the user/assistant merges in `repair_message_sequence`, or a context - # engine that rewrites in place). Such a row absorbed later durable rows whose ids are gone from - # `held`, so `max(held)` no longer names what the summary covered: capped there, those rows would be - # cloned live beside a summary that already contains them. A `here N` tail is exempt: it is copied - # verbatim (without the marker), so its ids are exact and, being newest, they lift the cap above - # anything a merged row earlier in `messages` absorbed. - if not any(isinstance(m, dict) for m in (verbatim_tail or ())) and not newest.get(_DB_PERSISTED_MARKER): + if agent._session_db.get_message_role(agent.session_id, newest_held) is None: return watermark - if agent._session_db.get_message_role(agent.session_id, max(held)) is None: - return watermark - return max(held) + return newest_held def _commit_compaction( diff --git a/agent/conversation_compression_manual.py b/agent/conversation_compression_manual.py index 16215b44c7..329a79d076 100644 --- a/agent/conversation_compression_manual.py +++ b/agent/conversation_compression_manual.py @@ -105,8 +105,15 @@ def compress_now( if skip_without_window and callable(has_content) and has_content(head) is False: return CompressResult("nothing_to_do", before, before, before_tokens, before_tokens, request) # An in-place commit archives every durable row under the lease watermark, the kept tail's included, so - # it must store the tail again itself. It gets copies because the insert writes row ids onto them. - tail_rows = [_fresh_compaction_message_copy(m) for m in tail] + # it must store the tail again itself. It gets copies because the insert writes row ids onto them. A copy + # keeps its source's row id only while the source is unchanged since it was loaded (marker present): a + # resume merge rewrote the row and dropped the marker, and its id no longer bounds what the row holds. + tail_rows = [] + for m in tail: + row = _fresh_compaction_message_copy(m) + if not m.get(_DB_PERSISTED_MARKER): + row.pop("_row_id", None) + tail_rows.append(row) try: compressed, _ = agent._compress_context( head, system_message, approx_tokens=before_tokens, focus_topic=request.focus_topic, force=True, diff --git a/tests/agent/test_conversation_compression_manual.py b/tests/agent/test_conversation_compression_manual.py index 3e81880cf4..e96a54eec4 100644 --- a/tests/agent/test_conversation_compression_manual.py +++ b/tests/agent/test_conversation_compression_manual.py @@ -269,18 +269,20 @@ def test_in_place_compress_never_leaves_two_live_copies_of_a_row(session_db, sta assert len(contents) == len(set(contents)) -def test_in_place_compress_never_clones_a_row_a_merge_already_carried(session_db): +@pytest.mark.parametrize("raw", ["", "here 2"]) +def test_in_place_compress_never_clones_a_row_a_merge_already_carried(session_db, raw): """Resume repair merges consecutive user rows into the first one's dict: that dict keeps its row id, drops the persisted marker, and the later row's id leaves the held history. Stopping the archive at the newest held id would then leave the later row above the stop, cloned as a concurrent append beside the merged row that - already carries its content. Two strings, so the exact-duplicate check above cannot see it.""" + already carries its content. Two strings, so the exact-duplicate check above cannot see it. With ``here 2`` + the merged row is the newest row of the verbatim tail, whose marker-swept copy must not vouch for its id.""" agent, _ = _stored_agent(session_db, _exchanges(10)) session_db.append_message("sid", "user", "first half of a split prompt") session_db.append_message("sid", "user", "second half 9931") held = session_db.get_resume_conversations("sid")[0] # the resume path merges the two user rows assert "second half 9931" in held[-1]["content"] and held[-1]["content"] != "second half 9931" - assert _compress(agent, held, "").status == "compressed" + assert _compress(agent, held, raw).status == "compressed" contents = [m["content"] for m in session_db.get_messages_as_conversation("sid")] assert sum("second half 9931" in c for c in contents) == 1 From a702791371a552b6d5e6a39197c5d34e0f31459d Mon Sep 17 00:00:00 2001 From: teknium1 <127238744+teknium1@users.noreply.github.com> Date: Thu, 24 Sep 2026 02:19:12 -0700 Subject: [PATCH 16/22] test(cli): stop leaking a real MCP discovery thread from the TUI launcher test tests/hermes_cli/test_tui_launcher_skips_plugin_discovery.py crashed CI with SIGABRT ("FATAL: exception not rethrown") after both tests passed (PR #120924, run 35967641887 attempt 1). Mechanism: the plain-chat case replaces sys.modules["hermes_cli.plugins"] with a SimpleNamespace spy that lacks has_enabled_agent_plugin_mcp. _prepare_agent_startup then calls mcp_startup.start_background_mcp_discovery, whose _has_configured_mcp_servers() probe imports that symbol from the stub, raises, and falls back to "assume configured" -- so a REAL cli-mcp-discovery daemon thread starts and is still importing tools.mcp_tool when pytest exits. At interpreter finalization the daemon thread re-takes the GIL, CPython ends it via pthread_exit, and glibc's forced unwind hits a non-rethrowing catch frame inside the C-extension import: abort. Repro (base, 40 runs under 20 busy CPUs): 2 SIGABRT (rc=134), thread alive at atexit in 38/40. Fixed: 0/40 aborts, 0/40 leaked threads. Fix: the spy also no-ops start_background_mcp_discovery -- plugin discovery is the only subject of this file. Assertions unchanged. --- .../test_tui_launcher_skips_plugin_discovery.py | 11 +++++++++++ 1 file changed, 11 insertions(+) diff --git a/tests/hermes_cli/test_tui_launcher_skips_plugin_discovery.py b/tests/hermes_cli/test_tui_launcher_skips_plugin_discovery.py index e029cbce85..935c10811d 100644 --- a/tests/hermes_cli/test_tui_launcher_skips_plugin_discovery.py +++ b/tests/hermes_cli/test_tui_launcher_skips_plugin_discovery.py @@ -14,6 +14,7 @@ import sys import types from hermes_cli import main as main_mod +from hermes_cli import mcp_startup def _install_discover_spy(monkeypatch): @@ -32,6 +33,16 @@ def _install_discover_spy(monkeypatch): start_background_plugin_discovery=_discover, ), ) + # The plain-chat path also arms MCP discovery. Its config probe imports + # ``hermes_cli.plugins`` (replaced by the stub above), fails, and falls + # back to "assume configured", which spawned a REAL ``cli-mcp-discovery`` + # daemon thread that was still importing ``tools.mcp_tool`` when pytest + # exited. A daemon thread inside a C-extension import at interpreter + # finalization dies via pthread_exit → glibc "FATAL: exception not + # rethrown" → SIGABRT. Plugin discovery is the only subject here. + monkeypatch.setattr( + mcp_startup, "start_background_mcp_discovery", lambda **_kw: None + ) return calls From 43d3d4e851479fab9e3bee7d8d209cd308fe782d Mon Sep 17 00:00:00 2001 From: brooklyn! Date: Thu, 24 Sep 2026 04:41:23 -0500 Subject: [PATCH 17/22] fix(journey): route memory edit/delete through MemoryStore._mutate MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A Journey memory edit or delete (Desktop panel, `hermes journey`, TUI `learning.*`) rewrote the whole MEMORY.md / USER.md from an unlocked snapshot via `_read_file` + `_write_file`: a memory the agent stored in between was dropped, and hand-edited content that does not round-trip through the § parser was reformatted with no .bak. Both mutations now go through `MemoryStore._mutate` — the memory tool's cross-process lock, re-read under lock and drift guard. The node id is resolved to its entry text inside the lock and matched by exact text against the re-read entries; a vanished target, drift or an unreadable file is refused with the store's own message instead of written over. Router/CLI/TUI response shapes are unchanged. Fixes #119668 --- agent/learning_mutations.py | 47 +++++++++----- tests/agent/test_learning_mutations.py | 89 ++++++++++++++++++++++++++ tools/memory_tool_store.py | 5 +- 3 files changed, 124 insertions(+), 17 deletions(-) diff --git a/agent/learning_mutations.py b/agent/learning_mutations.py index 2ee948af1b..05c2b69d61 100644 --- a/agent/learning_mutations.py +++ b/agent/learning_mutations.py @@ -5,7 +5,7 @@ Node ids (from ``agent.learning_graph``): skills → the skill name; memories for USER.md; ``index`` = position in the combined card list, MEMORY.md first). Shared by CLI ``hermes journey``, the TUI ``/journey`` overlay and the desktop. Deleting a skill *archives* it (``hermes curator restore`` recovers it); -deleting a memory rewrites its file. +deleting a memory rewrites its file under the memory tool's lock. """ from __future__ import annotations @@ -14,6 +14,7 @@ from pathlib import Path from typing import Any, Callable _MEMORY_FILES = {"memory": "MEMORY.md", "profile": "USER.md"} +_STORE_TARGETS = {"memory": "memory", "profile": "user"} # journey source -> MemoryStore target def parse_node_kind(node_id: str) -> str: @@ -35,7 +36,8 @@ def _locate_memory(node_id: str) -> tuple[Path, list[str], int]: """Resolve a memory node id to (file, all §-delimited entries, local index). Entries come from ``MemoryStore._read_file`` — the memory tool's own parser — so journey indices stay aligned with what the graph renders; a profile card's - local index is its global index minus the MEMORY.md card count.""" + local index is its global index minus the MEMORY.md card count. Read-only view: + mutations resolve the id again INSIDE ``_mutate_memory``'s lock.""" from hermes_constants import get_hermes_home from agent.learning_graph import _memory_cards from tools.memory_tool import MemoryStore @@ -56,11 +58,32 @@ def _locate_memory(node_id: str) -> tuple[Path, list[str], int]: return path, chunks, local -def _write_memory(path: Path, chunks: list[str]) -> None: - """Atomic temp-file + rename via the memory tool, so a concurrent reader - never sees a half-written file (and the §-join stays single-sourced).""" - from tools.memory_tool import MemoryStore - MemoryStore._write_file(path, [c.strip() for c in chunks if c.strip()]) +def _mutate_memory(node_id: str, replacement: str | None) -> dict[str, Any]: + """Replace (or, with ``replacement=None``, remove) the entry *node_id* names, through + ``MemoryStore._mutate`` — the memory tool's cross-process lock, re-read under lock and + drift guard (``.bak`` snapshot + refusal when the file wouldn't round-trip). The file is + shared with the live agent, so a read-modify-write from an unlocked snapshot silently + dropped whatever the agent stored in between and reformatted hand-edited files + (#119668). The id is resolved to its entry text INSIDE the lock and matched by exact + text against the store's re-read entries; a target gone under the lock is refused.""" + from tools.memory_tool import load_on_disk_store + + source, _ = _parse_memory_id(node_id) + name = _MEMORY_FILES[source] + message = f"deleted memory from {name}" if replacement is None else f"updated memory in {name}" + + def _apply(entries, _limit): + _, chunks, local = _locate_memory(node_id) + text = chunks[local].strip() + if text not in entries: + return {"success": False, "error": "memory node id is stale — refresh the graph"} + idx = entries.index(text) + return entries[:idx] + ([] if replacement is None else [replacement]) + entries[idx + 1:], message + + result = load_on_disk_store()._mutate(_STORE_TARGETS[source], _apply) + if not result.get("success"): + return {"ok": False, "message": result.get("error", f"{name} write failed")} + return {"ok": True, "message": message} def _clear_skill_cache() -> None: @@ -125,10 +148,7 @@ def _delete_skill(name: str) -> dict[str, Any]: def _delete_memory(node_id: str) -> dict[str, Any]: - path, chunks, local = _locate_memory(node_id) - del chunks[local] - _write_memory(path, chunks) - return {"ok": True, "message": f"deleted memory from {path.name}"} + return _mutate_memory(node_id, None) # ── Edit ──────────────────────────────────────────────────────────────────── @@ -151,7 +171,4 @@ def _edit_memory(node_id: str, content: str) -> dict[str, Any]: body = content.strip() if not body: return {"ok": False, "message": "empty memory — use delete to remove it"} - path, chunks, local = _locate_memory(node_id) - chunks[local] = body - _write_memory(path, chunks) - return {"ok": True, "message": f"updated memory in {path.name}"} + return _mutate_memory(node_id, body) diff --git a/tests/agent/test_learning_mutations.py b/tests/agent/test_learning_mutations.py index 7d2ef8d963..2927d6f58f 100644 --- a/tests/agent/test_learning_mutations.py +++ b/tests/agent/test_learning_mutations.py @@ -6,6 +6,8 @@ against a temp HERMES_HOME, never mocks — the id→file mapping is the whole p from __future__ import annotations +import threading + import pytest from agent import learning_mutations as lm @@ -91,3 +93,90 @@ def test_memory_writes_match_memory_tool_format(home): assert entries == ["alpha rewritten", "beta note"] assert path.read_text(encoding="utf-8") == ENTRY_DELIMITER.join(entries) + + +# ── Locking / drift (issue #119668) ───────────────────────────────────────── +# A Journey mutation shares MEMORY.md with the live agent's memory tool, so it must +# take the same lock, re-read under it and honour the drift guard — otherwise a +# memory the agent stored in the meantime is rewritten away from a stale snapshot. + + +def _race_memory_tool_add(monkeypatch, content: str) -> threading.Thread: + """Start a lock-respecting ``memory_tool`` writer that appends *content* the + moment the Journey mutation resolves its node id (the read half of its + read-modify-write), then give it a moment to land. Under a correct lock the + writer blocks until the mutation has written; without one it interleaves and + the mutation's write clobbers it.""" + from agent import learning_graph + from tools.memory_tool import MemoryStore + + located, landed = threading.Event(), threading.Event() + real_cards = learning_graph._memory_cards + + def _cards_then_let_writer_in(): + cards = real_cards() + located.set() + landed.wait(timeout=0.5) + return cards + + def _writer(): + located.wait(timeout=5) + MemoryStore().add("memory", content) + landed.set() + + monkeypatch.setattr(learning_graph, "_memory_cards", _cards_then_let_writer_in) + thread = threading.Thread(target=_writer, daemon=True) + thread.start() + return thread + + +def test_delete_memory_keeps_concurrent_memory_tool_add(home, monkeypatch): + from tools.memory_tool import MemoryStore + + writer = _race_memory_tool_add(monkeypatch, "gamma note") + assert lm.delete_node("memory:memory:0")["ok"] + writer.join(timeout=5) + + assert MemoryStore._read_file(home / "memories" / "MEMORY.md") == ["beta note", "gamma note"] + + +def test_edit_memory_keeps_concurrent_memory_tool_add(home, monkeypatch): + from tools.memory_tool import MemoryStore + + writer = _race_memory_tool_add(monkeypatch, "gamma note") + assert lm.edit_node("memory:memory:0", "alpha rewritten")["ok"] + writer.join(timeout=5) + + assert MemoryStore._read_file(home / "memories" / "MEMORY.md") == ["alpha rewritten", "beta note", "gamma note"] + + +@pytest.mark.parametrize("mutate", [lambda: lm.delete_node("memory:memory:0"), + lambda: lm.edit_node("memory:memory:0", "alpha rewritten")], + ids=["delete", "edit"]) +def test_memory_drift_is_refused_with_backup_like_memory_tool(home, mutate): + """Hand-edited content that wouldn't round-trip through the § parser is what + the memory tool's drift guard exists for: snapshot to .bak, refuse, leave the + file untouched. A Journey mutation must not reformat it silently.""" + from tools.memory_tool_store import _drift_error + + path = home / "memories" / "MEMORY.md" + raw = "alpha note\n§\n\n§\nbeta note\n" + path.write_text(raw, encoding="utf-8") + + res = mutate() + + assert not res["ok"] + assert path.read_text(encoding="utf-8") == raw + (backup,) = home.glob("memories/MEMORY.md.bak.*") + assert backup.read_text(encoding="utf-8") == raw + assert res["message"] == _drift_error(path, str(backup))["error"] + + +def test_unreadable_memory_file_is_refused_unchanged(home): + path = home / "memories" / "MEMORY.md" + path.write_bytes(b"alpha note\n\xc3\x28\n\xc2\xa7\nbeta") + + res = lm.delete_node("memory:memory:0") + + assert not res["ok"] and "could not be read" in res["message"] + assert path.read_bytes() == b"alpha note\n\xc3\x28\n\xc2\xa7\nbeta" diff --git a/tools/memory_tool_store.py b/tools/memory_tool_store.py index 2e29c83767..d340333ee1 100644 --- a/tools/memory_tool_store.py +++ b/tools/memory_tool_store.py @@ -462,8 +462,9 @@ class MemoryStore: @staticmethod def _write_file(path: Path, entries: List[str]): - """Atomic temp-file + rename: readers never see a truncated file. Also used by - agent/learning_mutations.py.""" + """Atomic temp-file + rename: readers never see a truncated file. Callers + hold ``_file_lock`` (via ``_mutate``): a bare write from an earlier snapshot + drops concurrent entries (#119668).""" try: atomic_write_text(path, ENTRY_DELIMITER.join(entries), tmp_prefix=".mem_") except OSError as e: From a88bef98b279651394e076e2987f50284ccdc2ff Mon Sep 17 00:00:00 2001 From: John Paul Soliva Date: Thu, 24 Sep 2026 09:26:00 +0900 Subject: [PATCH 18/22] fix(config): stop the migration ladder rewriting unversioned configs A config.yaml without _config_version reads as v0 and is exempt from the support floor, so the first `hermes update`, profile clone, `hermes doctor --fix` or docker boot ran every one-time migration step on it. Installers seed config.yaml from cli-config.yaml.example, which had no version, and targeted writers (`hermes config set`, /personality, the TUI/Desktop config writers) never stamp one, so this is the normal state of --skip-setup, non-TTY and Desktop (--non-interactive) installs. The value- and absence-based steps then reset the personality, raised the delegation caps, turned verify_on_stop off, shortened the curator windows, dropped model_catalog.ttl_hours and enabled plugins the user had installed but never enabled. - A config with no _config_version now gets only the steps keyed on a legacy key or identifier (LEGACY_KEY_STEPS), then the stamp. - cli-config.yaml.example carries _config_version, so every seeded config (install.sh, install.ps1, docker/stage2-hook.sh, doctor --fix) starts at the current schema. - docker_config_migrate.py no longer refuses a version-less volume with the "predates version 12" warning; like migrate_config() it migrates and stamps it. (cherry picked from commit 97ba11e07009f633662b0c7fa8701aa5b441bd22) --- cli-config.yaml.example | 5 ++ hermes_cli/config.py | 13 ++-- hermes_cli/config_migrations.py | 27 ++++++- scripts/docker_config_migrate.py | 7 +- .../test_config_unversioned_migration.py | 78 +++++++++++++++++++ tests/tools/test_docker_config_migrate.py | 17 ++-- 6 files changed, 126 insertions(+), 21 deletions(-) create mode 100644 tests/hermes_cli/test_config_unversioned_migration.py diff --git a/cli-config.yaml.example b/cli-config.yaml.example index 21ff9ea6e4..a622c100fd 100644 --- a/cli-config.yaml.example +++ b/cli-config.yaml.example @@ -4,6 +4,11 @@ # This file configures CLI behavior; only documented secret environment # variables in .env take precedence over their corresponding settings. +# Schema version of this file. The installers copy it to seed config.yaml, and +# `hermes update` uses it to know which one-time migrations the file already +# has. Hermes manages it: do not copy it into another config. +_config_version: 46 + # ============================================================================= # Database Configuration # ============================================================================= diff --git a/hermes_cli/config.py b/hermes_cli/config.py index ca846e01b7..9f6038fdd7 100644 --- a/hermes_cli/config.py +++ b/hermes_cli/config.py @@ -1310,17 +1310,14 @@ def migrate_config(interactive: bool = True, quiet: bool = False) -> Dict[str, A # Auto-migration support floor (v12): an EXPLICIT on-disk ``_config_version`` below the # floor is NOT migrated and NOT rewritten — surface a message and leave the file untouched - # (deep-merge supplies defaults at read time). A config with NO version key is a fresh - # minimal config, not an ancient install: it gets the normal ladder and a version stamp. + # (deep-merge supplies defaults at read time). A config with NO version key is not an + # ancient install: it gets only the legacy-key steps and a version stamp. # Missing/unparseable files never trip the floor gate. # Imported lazily because the steps call back into this module. from hermes_cli.config_migrations import ( - SUPPORT_FLOOR_VERSION, run_migrations, support_floor_message) + SUPPORT_FLOOR_VERSION, has_version_stamp, run_migrations, support_floor_message) - try: - has_explicit_version = "_config_version" in read_user_config_raw() - except Exception: - has_explicit_version = False + has_explicit_version = has_version_stamp() floor_refused = ( has_explicit_version and current_ver < SUPPORT_FLOOR_VERSION and current_ver < latest_ver) if floor_refused: @@ -1331,7 +1328,7 @@ def migrate_config(interactive: bool = True, quiet: bool = False) -> Dict[str, A if not quiet: print(f" ⚠ {msg}") else: - run_migrations(current_ver, results, quiet) + run_migrations(current_ver, results, quiet, unversioned=not has_explicit_version) _disable_suspicious_mcp_servers(results, quiet) _warn_invalid_platform_toolsets(results, quiet) diff --git a/hermes_cli/config_migrations.py b/hermes_cli/config_migrations.py index b011814ea2..a4f1eff099 100644 --- a/hermes_cli/config_migrations.py +++ b/hermes_cli/config_migrations.py @@ -36,6 +36,16 @@ def support_floor_message() -> str: "after reviewing the changelog.") +def has_version_stamp() -> bool: + """Whether config.yaml carries a ``_config_version`` key. ``check_config_version()`` reads a + missing one as 0, but such a file is never-stamped current-schema content, not an ancient + install: the floor must not refuse it and only :data:`LEGACY_KEY_STEPS` may run on it.""" + try: + return "_config_version" in _cfg().read_user_config_raw() + except Exception: + return False + + def _cfg(): """Return the live ``hermes_cli.config`` module (lazy, cycle-free, monkeypatch-friendly).""" from hermes_cli import config @@ -755,15 +765,26 @@ MIGRATIONS: Tuple[Tuple[int, Callable[[Dict[str, Any], bool], None]], ...] = ( (46, _migrate_to_46), ) +#: Steps triggered by a legacy key or identifier (a renamed or retired key, a removed plugin or +#: toolset, the plugin-era SOUL.md section): they carry its setting to where the runtime reads it +#: or drop what nothing reads, which is right however old the file is. A config.yaml with no +#: ``_config_version`` is current-schema content that was never stamped (installers seed it from +#: cli-config.yaml.example; targeted writers never stamp), so it gets only these: every other step +#: decides by a value or an absence that, in such a file, is the user's own choice. v13 is left +#: out: it clears OPENAI_MODEL from .env, a generic name Hermes never reads but the user's tools may. +LEGACY_KEY_STEPS = frozenset({12, 14, 16, 17, 29, 33, 38, 39, 41, 42, 43, 46}) -def run_migrations(current_ver: int, results: Dict[str, Any], quiet: bool) -> None: - """Apply every registered migration whose target version exceeds *current_ver*. + +def run_migrations( + current_ver: int, results: Dict[str, Any], quiet: bool, *, unversioned: bool = False) -> None: + """Apply every registered migration whose target version exceeds *current_ver*; a config + with no ``_config_version`` (*unversioned*) gets only :data:`LEGACY_KEY_STEPS`. *current_ver* is the on-disk schema version captured ONCE before any step runs and does not advance between steps — each step is gated on the same initial value. """ for target_ver, migration_fn in MIGRATIONS: - if current_ver < target_ver: + if current_ver < target_ver and (target_ver in LEGACY_KEY_STEPS or not unversioned): try: migration_fn(results, quiet) except Exception as exc: diff --git a/scripts/docker_config_migrate.py b/scripts/docker_config_migrate.py index 2e9d3d66b1..cede0d4b9f 100644 --- a/scripts/docker_config_migrate.py +++ b/scripts/docker_config_migrate.py @@ -16,6 +16,7 @@ from hermes_cli.config import ( from hermes_cli.config_backups import backup_config, list_config_backups from hermes_cli.config_migrations import ( SUPPORT_FLOOR_VERSION, + has_version_stamp, support_floor_message, ) from utils import env_var_enabled @@ -54,8 +55,10 @@ def main() -> int: # Below the auto-migration support floor: migrate_config() refuses (and # leaves the file untouched), so don't run the backup/verify dance that # would raise "did not advance config version" and block the boot. - # Warn-and-continue matches the CLI's fail-safe posture. - if current_ver < SUPPORT_FLOOR_VERSION: + # Warn-and-continue matches the CLI's fail-safe posture. A config with no + # _config_version (a volume seeded from the template) is not below the + # floor: migrate_config() stamps it. + if current_ver < SUPPORT_FLOOR_VERSION and has_version_stamp(): print( f"[config-migrate] WARNING: {support_floor_message()}", file=sys.stderr, diff --git a/tests/hermes_cli/test_config_unversioned_migration.py b/tests/hermes_cli/test_config_unversioned_migration.py new file mode 100644 index 0000000000..6b22b3d6d0 --- /dev/null +++ b/tests/hermes_cli/test_config_unversioned_migration.py @@ -0,0 +1,78 @@ +"""A config.yaml without ``_config_version`` is current-schema content that was never stamped: +the installers seed it from cli-config.yaml.example and targeted writers (``hermes config set``, +/personality) never stamp. The one-time migration ladder must not treat it as a v0 install — +its value- and absence-based steps would overwrite what the user chose.""" + +import shutil +from pathlib import Path + +import pytest +import yaml + +TEMPLATE = Path(__file__).resolve().parents[2] / "cli-config.yaml.example" + +USER_CHOICES = { + "delegation.max_concurrent_children": "3", + "delegation.max_iterations": "50", + "display.background_process_notifications": "all", + "agent.verify_on_stop": "true", + "curator.stale_after_days": "30", + "curator.archive_after_days": "90", + "model_catalog.ttl_hours": "24", +} +WATCHED = [*USER_CHOICES, "display.personality", "plugins.enabled"] + + +@pytest.fixture +def hermes_home(tmp_path, monkeypatch): + home = tmp_path / ".hermes" + home.mkdir() + monkeypatch.setattr(Path, "home", lambda: tmp_path) # sibling-profile scan stays in tmp + monkeypatch.setenv("HERMES_HOME", str(home)) + return home + + +def _raw(home: Path) -> dict: + return yaml.safe_load((home / "config.yaml").read_text(encoding="utf-8")) or {} + + +def _at(raw: dict, dotted: str): + for part in dotted.split("."): + raw = raw.get(part) if isinstance(raw, dict) else None + return raw + + +def test_update_keeps_user_values_of_an_unversioned_config_and_migrates_legacy_keys(hermes_home): + from hermes_cli.config import DEFAULT_CONFIG, set_config_value + from hermes_cli.personality import persist_personality + from hermes_cli.update_cmd import _check_and_apply_config_migration + + # Seeded before the template carried a stamp; early-2026 templates shipped this retired key. + (hermes_home / "config.yaml").write_text( + "compression:\n summary_model: google/gemini-3-flash-preview\n", encoding="utf-8") + # Installed with `hermes plugins install` and never enabled. + plugin = hermes_home / "plugins" / "notes-helper" + plugin.mkdir(parents=True) + (plugin / "plugin.yaml").write_text("name: notes-helper\nversion: 0.1.0\n", encoding="utf-8") + assert persist_personality("kawaii") + for key, value in USER_CHOICES.items(): + set_config_value(key, value) + before = _raw(hermes_home) + + _check_and_apply_config_migration() + + after = _raw(hermes_home) + assert {k: _at(after, k) for k in WATCHED} == {k: _at(before, k) for k in WATCHED} + assert "summary_model" not in after["compression"] + assert _at(after, "auxiliary.compression.model") == _at(before, "compression.summary_model") + assert after["_config_version"] == DEFAULT_CONFIG["_config_version"] + + +def test_config_seeded_from_the_template_reads_as_current(hermes_home): + """install.sh, install.ps1, docker/stage2-hook.sh and `hermes doctor --fix` copy the template.""" + from hermes_cli.config import check_config_version + + shutil.copy(TEMPLATE, hermes_home / "config.yaml") + + current, latest = check_config_version(raise_on_parse_error=True) + assert current == latest diff --git a/tests/tools/test_docker_config_migrate.py b/tests/tools/test_docker_config_migrate.py index ac9d3040be..dbdd4ef2b7 100644 --- a/tests/tools/test_docker_config_migrate.py +++ b/tests/tools/test_docker_config_migrate.py @@ -101,19 +101,20 @@ def test_docker_config_migrate_skips_below_floor_config_untouched(tmp_path: Path assert not list(tmp_path.glob("*.bak-*")) -def test_docker_config_migrate_skips_unversioned_config_untouched(tmp_path: Path) -> None: - """Unversioned configs coerce to version 0 — below the floor, so refused.""" +def test_docker_config_migrate_stamps_unversioned_config(tmp_path: Path) -> None: + """A config with no _config_version (a template-seeded volume) is not below the floor, as in + migrate_config(): it is stamped and keeps its values.""" config_path = tmp_path / "config.yaml" - original = yaml.safe_dump({"model": {"default": "m", "provider": "openrouter"}}) - config_path.write_text(original, encoding="utf-8") + model = {"default": "m", "provider": "openrouter"} + config_path.write_text(yaml.safe_dump({"model": model}), encoding="utf-8") proc = _run_migration(tmp_path) assert proc.returncode == 0, proc.stderr - assert "Migrating config schema" not in proc.stdout - assert "can no longer be auto-migrated" in proc.stderr - assert config_path.read_text(encoding="utf-8") == original - assert not list(tmp_path.glob("*.bak-*")) + assert "can no longer be auto-migrated" not in proc.stderr + raw = yaml.safe_load(config_path.read_text(encoding="utf-8")) + assert raw["_config_version"] == DEFAULT_CONFIG["_config_version"] + assert raw["model"] == model def test_docker_config_migrate_does_not_rewrite_invalid_yaml(tmp_path: Path) -> None: From 49aaa52d33d6aaa379ad372c08b2ebf4c7997a28 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 13:57:46 +0530 Subject: [PATCH 19/22] fix(config): keep v41 off the unversioned ladder and drop the stamp read guard v41 walks every profile SOUL.md and deletes any "## Messaging other agents" section on a bare heading match. A config.yaml with no _config_version says nothing about where that SOUL text came from, so an unstamped current config must not authorize rewriting a user-owned file; the versioned 40->41 path is unchanged. The unversioned regression test now seeds a user-authored section under that heading and asserts it survives. has_version_stamp() loses its bare except: migrate_config() already called check_config_version(raise_on_parse_error=True) and the docker script returns early on the (latest, latest) a parse failure yields, so the raw read cannot fail at this point. Co-authored-by: JoaoMarcos44 --- hermes_cli/config_migrations.py | 12 ++++++------ .../hermes_cli/test_config_unversioned_migration.py | 4 ++++ 2 files changed, 10 insertions(+), 6 deletions(-) diff --git a/hermes_cli/config_migrations.py b/hermes_cli/config_migrations.py index a4f1eff099..9eb21eb90e 100644 --- a/hermes_cli/config_migrations.py +++ b/hermes_cli/config_migrations.py @@ -39,11 +39,9 @@ def support_floor_message() -> str: def has_version_stamp() -> bool: """Whether config.yaml carries a ``_config_version`` key. ``check_config_version()`` reads a missing one as 0, but such a file is never-stamped current-schema content, not an ancient - install: the floor must not refuse it and only :data:`LEGACY_KEY_STEPS` may run on it.""" - try: - return "_config_version" in _cfg().read_user_config_raw() - except Exception: - return False + install: the floor must not refuse it and only :data:`LEGACY_KEY_STEPS` may run on it. + Callers have already gone through ``check_config_version()``, so the read cannot fail here.""" + return "_config_version" in _cfg().read_user_config_raw() def _cfg(): @@ -772,7 +770,9 @@ MIGRATIONS: Tuple[Tuple[int, Callable[[Dict[str, Any], bool], None]], ...] = ( #: cli-config.yaml.example; targeted writers never stamp), so it gets only these: every other step #: decides by a value or an absence that, in such a file, is the user's own choice. v13 is left #: out: it clears OPENAI_MODEL from .env, a generic name Hermes never reads but the user's tools may. -LEGACY_KEY_STEPS = frozenset({12, 14, 16, 17, 29, 33, 38, 39, 41, 42, 43, 46}) +#: v41 is left out too: it rewrites profile SOUL.md on a heading match, an artifact whose +#: provenance the config stamp says nothing about. +LEGACY_KEY_STEPS = frozenset({12, 14, 16, 17, 29, 33, 38, 39, 42, 43, 46}) def run_migrations( diff --git a/tests/hermes_cli/test_config_unversioned_migration.py b/tests/hermes_cli/test_config_unversioned_migration.py index 6b22b3d6d0..b3ecc8fa59 100644 --- a/tests/hermes_cli/test_config_unversioned_migration.py +++ b/tests/hermes_cli/test_config_unversioned_migration.py @@ -54,6 +54,9 @@ def test_update_keeps_user_values_of_an_unversioned_config_and_migrates_legacy_k plugin = hermes_home / "plugins" / "notes-helper" plugin.mkdir(parents=True) (plugin / "plugin.yaml").write_text("name: notes-helper\nversion: 0.1.0\n", encoding="utf-8") + # User-authored SOUL.md section that happens to carry the v41 heading. + soul = "# Me\n\n## Messaging other agents\nmy own notes\n\n## Prefs\nkeep\n" + (hermes_home / "SOUL.md").write_text(soul, encoding="utf-8") assert persist_personality("kawaii") for key, value in USER_CHOICES.items(): set_config_value(key, value) @@ -66,6 +69,7 @@ def test_update_keeps_user_values_of_an_unversioned_config_and_migrates_legacy_k assert "summary_model" not in after["compression"] assert _at(after, "auxiliary.compression.model") == _at(before, "compression.summary_model") assert after["_config_version"] == DEFAULT_CONFIG["_config_version"] + assert (hermes_home / "SOUL.md").read_text(encoding="utf-8") == soul def test_config_seeded_from_the_template_reads_as_current(hermes_home): From 92324790f2ee683a4f0c46fbba1b012c20ee768a Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 14:16:52 +0530 Subject: [PATCH 20/22] refactor(config): read the version stamp once and drop has_version_stamp migrate_config() and the docker boot script each parsed config.yaml twice: check_config_version() coerced a missing `_config_version` to 0 and threw the "was it present" bit away, so has_version_stamp() re-read the file to recover it, guarded only by a call-order promise in its docstring. That promise did not hold for the docker script, which used the tolerant check: a list-rooted config.yaml read as "unversioned, not below the floor", ran the backup + migrate_config() dance and exited 1 (base: floor warning, exit 0). Factor the read into _read_config_version_stamp() -> (Optional[int], latest); None means the mapping has no stamp. check_config_version() is a thin wrapper (None -> 0) so its 10 callers see identical output. migrate_config() and the docker script decide `unversioned` from that single read; the docker script now does the strict read itself and leaves an unparseable or non-mapping file alone with a warning and exit 0, matching its invalid-YAML posture. has_version_stamp() is deleted. Docker tests that mocked the pre-check now mock the new helper. --- hermes_cli/config.py | 35 ++++++++++++++++------- hermes_cli/config_migrations.py | 8 ------ scripts/docker_config_migrate.py | 19 ++++++++---- tests/tools/test_docker_config_migrate.py | 10 +++++-- 4 files changed, 45 insertions(+), 27 deletions(-) diff --git a/hermes_cli/config.py b/hermes_cli/config.py index 9f6038fdd7..5450a88314 100644 --- a/hermes_cli/config.py +++ b/hermes_cli/config.py @@ -910,14 +910,12 @@ def _coerce_config_version(value: Any) -> int: return max(version, 0) -def check_config_version(*, raise_on_parse_error: bool = False) -> Tuple[int, int]: - """Return ``(current_version, latest_version)`` from the raw on-disk config. - Reads the raw file rather than ``load_config()``: the deep-merge would make a file lacking - ``_config_version`` inherit the latest version, hiding that the schema was never migrated. - Invalid YAML gets a parse warning, not an automatic schema rewrite. Tolerant runtime status - callers keep the historical latest/latest fallback for malformed YAML; mutation and explicit - validation paths set ``raise_on_parse_error`` so a parse failure or a non-mapping root cannot - be mistaken for an up-to-date config.""" +def _read_config_version_stamp(*, raise_on_parse_error: bool = False) -> Tuple[Optional[int], int]: + """Single raw read behind ``check_config_version()``: ``(stamp, latest_version)`` where + *stamp* is ``None`` when config.yaml parsed but carries no ``_config_version`` key (a + never-stamped current-schema file, not an ancient install — ``migrate_config()`` gives it only + the legacy-key steps). A missing file, or malformed YAML under a tolerant caller, reads as + ``latest`` exactly as ``check_config_version()`` always reported it.""" latest = _coerce_config_version(DEFAULT_CONFIG.get("_config_version", 1)) or 1 config_path = get_config_path() if not config_path.exists(): @@ -945,9 +943,23 @@ def check_config_version(*, raise_on_parse_error: bool = False) -> Tuple[int, in f"a mapping, got {type(config).__name__}" ) config = {} + if "_config_version" not in config: + return None, latest return _coerce_config_version(config.get("_config_version")), latest +def check_config_version(*, raise_on_parse_error: bool = False) -> Tuple[int, int]: + """Return ``(current_version, latest_version)`` from the raw on-disk config. + Reads the raw file rather than ``load_config()``: the deep-merge would make a file lacking + ``_config_version`` inherit the latest version, hiding that the schema was never migrated. + Invalid YAML gets a parse warning, not an automatic schema rewrite. Tolerant runtime status + callers keep the historical latest/latest fallback for malformed YAML; mutation and explicit + validation paths set ``raise_on_parse_error`` so a parse failure or a non-mapping root cannot + be mistaken for an up-to-date config. A file with no version key reads as 0.""" + stamp, latest = _read_config_version_stamp(raise_on_parse_error=raise_on_parse_error) + return (0 if stamp is None else stamp), latest + + # ---- Config structure validation ---- # DEFAULT_CONFIG is the single source of truth for documented roots; the set is derived so new @@ -1299,7 +1311,8 @@ def migrate_config(interactive: bool = True, quiet: bool = False) -> Dict[str, A # Validate config.yaml before any migration side effect: sanitize_env_file() rewrites .env, # which must not happen when the migration will be refused for malformed YAML. - current_ver, latest_ver = check_config_version(raise_on_parse_error=True) + stamp, latest_ver = _read_config_version_stamp(raise_on_parse_error=True) + current_ver = 0 if stamp is None else stamp try: fixes = sanitize_env_file() @@ -1315,9 +1328,9 @@ def migrate_config(interactive: bool = True, quiet: bool = False) -> Dict[str, A # Missing/unparseable files never trip the floor gate. # Imported lazily because the steps call back into this module. from hermes_cli.config_migrations import ( - SUPPORT_FLOOR_VERSION, has_version_stamp, run_migrations, support_floor_message) + SUPPORT_FLOOR_VERSION, run_migrations, support_floor_message) - has_explicit_version = has_version_stamp() + has_explicit_version = stamp is not None floor_refused = ( has_explicit_version and current_ver < SUPPORT_FLOOR_VERSION and current_ver < latest_ver) if floor_refused: diff --git a/hermes_cli/config_migrations.py b/hermes_cli/config_migrations.py index 9eb21eb90e..41766a71d8 100644 --- a/hermes_cli/config_migrations.py +++ b/hermes_cli/config_migrations.py @@ -36,14 +36,6 @@ def support_floor_message() -> str: "after reviewing the changelog.") -def has_version_stamp() -> bool: - """Whether config.yaml carries a ``_config_version`` key. ``check_config_version()`` reads a - missing one as 0, but such a file is never-stamped current-schema content, not an ancient - install: the floor must not refuse it and only :data:`LEGACY_KEY_STEPS` may run on it. - Callers have already gone through ``check_config_version()``, so the read cannot fail here.""" - return "_config_version" in _cfg().read_user_config_raw() - - def _cfg(): """Return the live ``hermes_cli.config`` module (lazy, cycle-free, monkeypatch-friendly).""" from hermes_cli import config diff --git a/scripts/docker_config_migrate.py b/scripts/docker_config_migrate.py index cede0d4b9f..9a2237f5fb 100644 --- a/scripts/docker_config_migrate.py +++ b/scripts/docker_config_migrate.py @@ -8,6 +8,8 @@ from pathlib import Path from typing import Iterable from hermes_cli.config import ( + InvalidUserConfigError, + _read_config_version_stamp, check_config_version, get_config_path, get_env_path, @@ -16,7 +18,6 @@ from hermes_cli.config import ( from hermes_cli.config_backups import backup_config, list_config_backups from hermes_cli.config_migrations import ( SUPPORT_FLOOR_VERSION, - has_version_stamp, support_floor_message, ) from utils import env_var_enabled @@ -48,7 +49,15 @@ def main() -> int: print("[config-migrate] HERMES_SKIP_CONFIG_MIGRATION is set; skipping config migration") return 0 - current_ver, latest_ver = check_config_version() + # Strict read: malformed YAML or a non-mapping root is left alone with a warning (the + # tolerant check_config_version() already printed one) and the boot continues, instead of + # running the backup/migrate dance that migrate_config() would refuse anyway. + try: + stamp, latest_ver = _read_config_version_stamp(raise_on_parse_error=True) + except InvalidUserConfigError as exc: + print(f"[config-migrate] WARNING: {exc}; leaving config.yaml untouched", file=sys.stderr) + return 0 + current_ver = 0 if stamp is None else stamp if current_ver >= latest_ver: return 0 @@ -56,9 +65,9 @@ def main() -> int: # leaves the file untouched), so don't run the backup/verify dance that # would raise "did not advance config version" and block the boot. # Warn-and-continue matches the CLI's fail-safe posture. A config with no - # _config_version (a volume seeded from the template) is not below the - # floor: migrate_config() stamps it. - if current_ver < SUPPORT_FLOOR_VERSION and has_version_stamp(): + # _config_version (stamp None: a volume seeded from the template) is not + # below the floor: migrate_config() stamps it. + if stamp is not None and current_ver < SUPPORT_FLOOR_VERSION: print( f"[config-migrate] WARNING: {support_floor_message()}", file=sys.stderr, diff --git a/tests/tools/test_docker_config_migrate.py b/tests/tools/test_docker_config_migrate.py index dbdd4ef2b7..ee3cc4cb48 100644 --- a/tests/tools/test_docker_config_migrate.py +++ b/tests/tools/test_docker_config_migrate.py @@ -155,7 +155,9 @@ def test_docker_config_migrate_restores_backups_after_failed_migration( config_path.write_text(original_config, encoding="utf-8") env_path.write_text(original_env, encoding="utf-8") - monkeypatch.setattr(module, "check_config_version", lambda: (12, DEFAULT_CONFIG["_config_version"])) + monkeypatch.setattr( + module, "_read_config_version_stamp", + lambda *, raise_on_parse_error=False: (12, DEFAULT_CONFIG["_config_version"])) monkeypatch.setattr(module, "get_config_path", lambda: config_path) monkeypatch.setattr(module, "get_env_path", lambda: env_path) @@ -186,8 +188,10 @@ def test_docker_config_migrate_restores_backups_when_version_does_not_advance( config_path.write_text(original_config, encoding="utf-8") env_path.write_text(original_env, encoding="utf-8") - calls = iter([(12, DEFAULT_CONFIG["_config_version"]), (12, DEFAULT_CONFIG["_config_version"])]) - monkeypatch.setattr(module, "check_config_version", lambda: next(calls)) + monkeypatch.setattr( + module, "_read_config_version_stamp", + lambda *, raise_on_parse_error=False: (12, DEFAULT_CONFIG["_config_version"])) + monkeypatch.setattr(module, "check_config_version", lambda: (12, DEFAULT_CONFIG["_config_version"])) monkeypatch.setattr(module, "get_config_path", lambda: config_path) monkeypatch.setattr(module, "get_env_path", lambda: env_path) From 19bd20fcb6b58312d4ef7b38b5adf064bd4b01d4 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 14:17:36 +0530 Subject: [PATCH 21/22] test(config): tie LEGACY_KEY_STEPS to the migration ladder LEGACY_KEY_STEPS is a frozenset of bare version ints beside MIGRATIONS with nothing checking they agree: a typo or a retired step would silently change what an unversioned config.yaml receives. Assert in the existing registry test that every allowlisted version names a ladder step, and point authors of new steps at the allowlist from the MIGRATIONS header. --- hermes_cli/config_migrations.py | 3 ++- tests/hermes_cli/test_config.py | 4 ++++ 2 files changed, 6 insertions(+), 1 deletion(-) diff --git a/hermes_cli/config_migrations.py b/hermes_cli/config_migrations.py index 41766a71d8..94e94b8042 100644 --- a/hermes_cli/config_migrations.py +++ b/hermes_cli/config_migrations.py @@ -636,7 +636,8 @@ def _migrate_to_46(results: Dict[str, Any], quiet: bool) -> None: #: earlier steps' writes via read_raw_config() (filesystem state). v12 is the support floor: #: configs already AT v12 still get every step below; only configs BELOW 12 are refused by the #: floor gate in run_migrations()'s caller. Versions absent here (15, 18-20, 22, 24, 26-28, 30) -#: only added a schema default that runtime merging supplies without a write. +#: only added a schema default that runtime merging supplies without a write. When adding a step, +#: decide whether it belongs in LEGACY_KEY_STEPS below (the only steps an unversioned file gets). MIGRATIONS: Tuple[Tuple[int, Callable[[Dict[str, Any], bool], None]], ...] = ( (12, _migrate_to_12), (13, _migrate_to_13), diff --git a/tests/hermes_cli/test_config.py b/tests/hermes_cli/test_config.py index 601bb58296..6df6563746 100644 --- a/tests/hermes_cli/test_config.py +++ b/tests/hermes_cli/test_config.py @@ -858,6 +858,7 @@ class TestConfigSupportFloor: def test_registry_has_no_targets_below_floor(self): from hermes_cli.config_migrations import ( + LEGACY_KEY_STEPS, MIGRATIONS, SUPPORT_FLOOR_VERSION, ) @@ -867,6 +868,9 @@ class TestConfigSupportFloor: # v12's own step is retained: a config AT v11 is refused, but a # config AT v12 must still receive every remaining migration. assert MIGRATIONS[0][0] == 12 + # The unversioned allowlist is a parallel registry of ints: every entry must name a + # step that exists on the ladder, or an unversioned config silently drifts. + assert LEGACY_KEY_STEPS <= {target for target, _ in MIGRATIONS} # ── Parity fixtures ────────────────────────────────────────────── # Expected outputs captured by running migrate_config from origin/main From 03584e41b992747e1b6fe14f0766fde05a838306 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 14:29:10 +0530 Subject: [PATCH 22/22] test(docker): pin the non-mapping config.yaml boot path The strict version read in docker_config_migrate is what keeps a list-root config.yaml on the warn-and-continue path; only a probe covered it. Fold the list-root case into the existing invalid-YAML test (reverting to the tolerant read now goes red) and correct the comment about where the warning comes from. --- scripts/docker_config_migrate.py | 6 +++--- tests/tools/test_docker_config_migrate.py | 6 +++--- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/scripts/docker_config_migrate.py b/scripts/docker_config_migrate.py index 9a2237f5fb..bd111f0d0f 100644 --- a/scripts/docker_config_migrate.py +++ b/scripts/docker_config_migrate.py @@ -49,9 +49,9 @@ def main() -> int: print("[config-migrate] HERMES_SKIP_CONFIG_MIGRATION is set; skipping config migration") return 0 - # Strict read: malformed YAML or a non-mapping root is left alone with a warning (the - # tolerant check_config_version() already printed one) and the boot continues, instead of - # running the backup/migrate dance that migrate_config() would refuse anyway. + # Strict read: malformed YAML or a non-mapping root is left alone with a warning and the + # boot continues, instead of running the backup/migrate dance that migrate_config() would + # refuse anyway. try: stamp, latest_ver = _read_config_version_stamp(raise_on_parse_error=True) except InvalidUserConfigError as exc: diff --git a/tests/tools/test_docker_config_migrate.py b/tests/tools/test_docker_config_migrate.py index ee3cc4cb48..8bb1cad6d8 100644 --- a/tests/tools/test_docker_config_migrate.py +++ b/tests/tools/test_docker_config_migrate.py @@ -117,16 +117,16 @@ def test_docker_config_migrate_stamps_unversioned_config(tmp_path: Path) -> None assert raw["model"] == model -def test_docker_config_migrate_does_not_rewrite_invalid_yaml(tmp_path: Path) -> None: +@pytest.mark.parametrize("original", ["model: [unterminated\n", "- a\n- b\n"]) +def test_docker_config_migrate_does_not_rewrite_invalid_yaml(tmp_path: Path, original: str) -> None: config_path = tmp_path / "config.yaml" - original = "model: [unterminated\n" config_path.write_text(original, encoding="utf-8") proc = _run_migration(tmp_path) assert proc.returncode == 0, proc.stderr assert "Migrating config schema" not in proc.stdout - assert "hermes config:" in proc.stderr + assert "leaving config.yaml untouched" in proc.stderr assert config_path.read_text(encoding="utf-8") == original assert not list(tmp_path.glob("*.bak-*"))