diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 34c0eac45c..209907df7b 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -1250,12 +1250,55 @@ 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 result[0] is messages + + 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: + 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, + ) + return messages, _resolve_fallback_prompt() + try: settled, result = _await_worker_within_budget( future, fence, idle=idle, ceiling=ceiling, wait_started=wait_started ) 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. + 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) + return _recover_from_stall() return result # F6: a not-yet-started future must not linger as a stale queued job. @@ -1282,39 +1325,22 @@ 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. + if stall_fallback and _is_unchanged_snapshot(result): + 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) - 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() finally: if not handled_exit: # Any unwind while waiting: revoke commit admission and release the worker's @@ -3581,6 +3607,48 @@ 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`` 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 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 + + 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 + if agent._session_db.get_message_role(agent.session_id, newest_held) is None: + return watermark + return newest_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, @@ -3647,8 +3715,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/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/agent/learning_mutations.py b/agent/learning_mutations.py index 2c5016e667..dafa65a887 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,14 +58,32 @@ def _locate_memory(node_id: str) -> tuple[Path, list[str], int]: return path, chunks, local -# ── Helpers ───────────────────────────────────────────────────────────────── +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 _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 _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: @@ -128,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 ──────────────────────────────────────────────────────────────────── @@ -154,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/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/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 diff --git a/hermes_cli/config.py b/hermes_cli/config.py index 23cf21302e..61be52a827 100644 --- a/hermes_cli/config.py +++ b/hermes_cli/config.py @@ -926,14 +926,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(): @@ -961,9 +959,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 @@ -1315,7 +1327,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() @@ -1326,17 +1339,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) - try: - has_explicit_version = "_config_version" in read_user_config_raw() - except Exception: - has_explicit_version = False + 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: @@ -1347,7 +1357,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 6847cea1b7..04d4b157fc 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), @@ -755,15 +756,28 @@ 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. +#: 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(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/hermes_cli/web_routers/_common.py b/hermes_cli/web_routers/_common.py index 2399f90e24..8ccc851465 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 @@ -118,6 +119,36 @@ 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))}»" + + +# 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 "") + # 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 + + # 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 606495eca3..b5001978a9 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, _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 @@ -246,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", ""), @@ -292,6 +295,10 @@ 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 @@ -360,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}}}" @@ -620,6 +627,10 @@ 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: + # ``${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 _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 entry.pop("api_key", None) diff --git a/hermes_cli/web_routers/messaging.py b/hermes_cli/web_routers/messaging.py index d4722c3ea2..be5ab2eb15 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, @@ -232,7 +235,7 @@ def _messaging_platform_payload( 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"]) ] @@ -886,16 +889,25 @@ async def update_messaging_platform(platform_id: str, body: MessagingPlatformUpd def _apply(): with _profile_scope(target_profile): + updates: dict[str, str] = {} + + # Validate the whole request 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 + 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 + + 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) diff --git a/scripts/docker_config_migrate.py b/scripts/docker_config_migrate.py index 2e9d3d66b1..bd111f0d0f 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, @@ -47,15 +49,25 @@ 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 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 # 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 (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/agent/test_compression_attempt_lifecycle.py b/tests/agent/test_compression_attempt_lifecycle.py index ef7022db3b..52b7c27142 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,15 @@ 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 +150,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=0.3, + 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: diff --git a/tests/agent/test_conversation_compression_manual.py b/tests/agent/test_conversation_compression_manual.py index 1808302305..e96a54eec4 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,77 @@ 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", "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; 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") + 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)) + + +@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. 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, 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 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/tests/hermes_cli/test_config.py b/tests/hermes_cli/test_config.py index ed0304ca55..0ccc94d509 100644 --- a/tests/hermes_cli/test_config.py +++ b/tests/hermes_cli/test_config.py @@ -848,6 +848,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, ) @@ -857,6 +858,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 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..5f2b777a49 --- /dev/null +++ b/tests/hermes_cli/test_config_unversioned_migration.py @@ -0,0 +1,82 @@ +"""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 hermes_yaml as 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") + # 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) + 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"] + assert (hermes_home / "SOUL.md").read_text(encoding="utf-8") == soul + + +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/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 diff --git a/tests/hermes_cli/test_web_server.py b/tests/hermes_cli/test_web_server.py index 4b95ec8157..0d2dd4ccb8 100644 --- a/tests/hermes_cli/test_web_server.py +++ b/tests/hermes_cli/test_web_server.py @@ -2133,6 +2133,73 @@ 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): + """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" + real = "sk-live-secret-abcdef1234567890" + 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()[key] == real + + rotated = "sk-rotated-secret-0987654321" + save_env_value(key, rotated) + for stale in (preview, redact_key(real)): + response = self.client.put("/api/env", json={"key": key, "value": stale}) + assert response.status_code == 400 + assert load_env()[key] == rotated + + 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) + response = self.client.put( + "/api/messaging/platforms/discord", + 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 diff --git a/tests/tools/test_docker_config_migrate.py b/tests/tools/test_docker_config_migrate.py index ff80e86ba8..f5ccd92ca9 100644 --- a/tests/tools/test_docker_config_migrate.py +++ b/tests/tools/test_docker_config_migrate.py @@ -101,31 +101,32 @@ 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" + 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 "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 + + +@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 = yaml.safe_dump({"model": {"default": "m", "provider": "openrouter"}}) 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 "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-*")) - - -def test_docker_config_migrate_does_not_rewrite_invalid_yaml(tmp_path: Path) -> 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-*")) @@ -154,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) @@ -185,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) 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: