From 3cd4d6d69a71187eb5ae88daa9ca879b0fb9d443 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 11:23:17 -0700 Subject: [PATCH] =?UTF-8?q?refactor(plugins/memory):=20hindsight=20?= =?UTF-8?q?=E2=80=94=20split=20settings/embedded/setup=20modules,=20unify?= =?UTF-8?q?=20config=20parsing,=20remove=20dead=20helpers?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- plugins/memory/hindsight/__init__.py | 2427 +++++------------ plugins/memory/hindsight/embedded.py | 195 ++ plugins/memory/hindsight/settings.py | 158 ++ plugins/memory/hindsight/setup.py | 215 ++ .../memory/test_hindsight_env_perms.py | 4 +- .../plugins/memory/test_hindsight_provider.py | 4 +- 6 files changed, 1289 insertions(+), 1714 deletions(-) create mode 100644 plugins/memory/hindsight/embedded.py create mode 100644 plugins/memory/hindsight/settings.py create mode 100644 plugins/memory/hindsight/setup.py diff --git a/plugins/memory/hindsight/__init__.py b/plugins/memory/hindsight/__init__.py index cf3b9f618f..27d9910fbf 100644 --- a/plugins/memory/hindsight/__init__.py +++ b/plugins/memory/hindsight/__init__.py @@ -3,12 +3,6 @@ Long-term memory with knowledge graph, entity resolution, and multi-strategy retrieval. Supports cloud (API key) and local modes. -Configurable request timeout via HINDSIGHT_TIMEOUT env var or config.json. -Configurable embedded daemon idle timeout via HINDSIGHT_IDLE_TIMEOUT env var -or config.json idle_timeout. - -Original PR #1811 by benfrank241, adapted to MemoryProvider ABC. - Config via environment variables: HINDSIGHT_API_KEY — API key for Hindsight Cloud HINDSIGHT_BANK_ID — memory bank identifier (default: hermes) @@ -33,7 +27,6 @@ from __future__ import annotations import asyncio import atexit import contextvars -import importlib import json import logging import os @@ -41,169 +34,69 @@ import queue import sys import threading import time - from dataclasses import dataclass from datetime import datetime, timezone +from pathlib import Path from typing import Any, Callable, Dict, List, Optional -from agent.secret_scope import get_secret - from agent.memory_provider import MemoryProvider, RecallStatus +from agent.secret_scope import get_secret +from hermes_cli.config import cfg_get from hermes_constants import get_hermes_home from hermes_time import now as _hermes_now from tools.registry import tool_error -from hermes_cli.config import cfg_get + +from .embedded import ( # noqa: F401 (re-exported; tests patch/import via this module) + _PORT_HEALTH_GRACE_ENV, + _RETRIABLE_CONNECTION_MARKERS, + _build_embedded_profile_env, + _check_local_runtime, + _embedded_profile_env_path, + _export_port_health_grace_timeout, + _load_simple_env, + _local_runtime_hint, + _materialize_embedded_profile_env, + _secure_write_profile_env, + _validate_profile_env_permissions, +) +from .settings import ( # noqa: F401 + _DEFAULT_API_URL, + _DEFAULT_IDLE_TIMEOUT, + _DEFAULT_LOCAL_URL, + _DEFAULT_RETAIN_SOURCE, + _DEFAULT_TIMEOUT, + _HINDSIGHT_GLYPH, + _MIN_CLIENT_VERSION, + _MIN_VERSION_FOR_UPDATE_MODE_APPEND, + _OBSERVATION_SCOPE_KEYWORDS, + _PROVIDER_DEFAULT_MODELS, + _VALID_BUDGETS, + _daemon_llm_provider, + _normalize_observation_scopes, + _normalize_retain_tags, + _parse_int_setting, + _resolve_bank_id_template, + _sanitize_bank_segment, +) logger = logging.getLogger(__name__) +_LOCAL_MODES = {"local", "local_embedded"} +_RETAIN_CONTEXT_DEFAULT = "conversation between Hermes Agent and the User" + @dataclass(frozen=True) class _RecallResult: - """Text + memory count from one recall. - - Carrying the count alongside the text lets the deterministic recall - indicator report "recalled N memories" accurately without re-parsing the - formatted bullet list. ``count`` is 0 for a reflect synthesis (no discrete - memories) or on error. - """ + """Text + memory count from one recall, so the deterministic recall indicator + can report "recalled N memories" without re-parsing the formatted text. + ``count`` is 0 for a reflect synthesis or on error.""" text: str count: int -_DEFAULT_API_URL = "https://api.hindsight.vectorize.io" -_DEFAULT_LOCAL_URL = "http://localhost:8888" -# Keep in sync with tools/lazy_deps.py ("memory.hindsight") and plugin.yaml. -_MIN_CLIENT_VERSION = "0.6.1" -_DEFAULT_TIMEOUT = 120 # seconds — cloud API can take 30-40s per request -_DEFAULT_IDLE_TIMEOUT = 300 # seconds — Hindsight embedded daemon default -# ``metadata.source`` stamped on retained memories — OPT-IN, empty by default. -# AGENTS.md forbids shipping third-party attribution tags on-by-default until a -# generic user-facing opt-in exists, so this stays unset unless the user sets it -# via the ``retain_source`` config key or HINDSIGHT_RETAIN_SOURCE (e.g. "hermes"). -_DEFAULT_RETAIN_SOURCE = "" -# Hindsight brand mark — the logo is an eye ringed by graph nodes. Used for -# the deterministic recall/retain indicators (overrides the generic core default). -_HINDSIGHT_GLYPH = "👁️" -# Mirrors hindsight-integrations/openclaw — Hindsight 0.5.0 added -# `update_mode='append'` semantics on retain (vectorize-io/hindsight#932). -# Without it, reusing a stable session-scoped document_id silently -# overwrites prior turns server-side, so we keep the per-process -# unique document_id fallback for older APIs. -_MIN_VERSION_FOR_UPDATE_MODE_APPEND = "0.5.0" -_VALID_BUDGETS = {"low", "mid", "high"} -_PROVIDER_DEFAULT_MODELS = { - "openai": "gpt-4o-mini", - "anthropic": "claude-haiku-4-5", - "gemini": "gemini-3.6-flash", - "groq": "openai/gpt-oss-120b", - "openrouter": "qwen/qwen3.5-9b", - "minimax": "MiniMax-M2.7", - "ollama": "gemma3:12b", - "lmstudio": "local-model", - "openai_compatible": "your-model-name", -} - -def _parse_int_setting(value: Any, default: int) -> int: - """Parse an integer config/env value, falling back on invalid input.""" - if value is None or value == "": - return default - try: - return int(value) - except (TypeError, ValueError): - logger.warning("Invalid integer Hindsight setting %r; using default %s", value, default) - return default - - -# Env var the embedded daemon manager reads (at import time, as a module-level -# constant) to size the grace window it waits for a slow /health before -# declaring a daemon stale and killing it. Default upstream is 30s; on -# resource-contended hosts a busy daemon can exceed a single 2s health check -# and get needlessly killed + restarted (issue #13125 comment thread). We -# surface it as plugin config so users can raise it without hand-setting an -# env var, consistent with "config.json, not raw env vars". -_PORT_HEALTH_GRACE_ENV = "HINDSIGHT_EMBED_PORT_HEALTH_GRACE_TIMEOUT" - - -def _export_port_health_grace_timeout(config: dict[str, Any]) -> None: - """Export the embedded-daemon health grace timeout to the process env. - - Must run BEFORE ``hindsight_embed.daemon_embed_manager`` is imported, - because the package reads the env var into a module-level constant at - import time. We only set it when the user configured a value AND the - env var isn't already set, so an explicit env override always wins. - """ - raw = config.get("port_health_grace_timeout") - if raw is None or raw == "": - return - try: - seconds = float(raw) - except (TypeError, ValueError): - logger.warning( - "Invalid Hindsight port_health_grace_timeout %r; ignoring.", raw - ) - return - if seconds < 0: - logger.warning( - "Negative Hindsight port_health_grace_timeout %r; ignoring.", raw - ) - return - # setdefault: an explicit env var the operator set wins over config. - os.environ.setdefault(_PORT_HEALTH_GRACE_ENV, repr(seconds)) - - -def _check_local_runtime() -> tuple[bool, str | None]: - """Return whether local embedded Hindsight imports cleanly. - - On older CPUs, importing the local Hindsight stack can raise a runtime - error from NumPy before the daemon starts. Treat that as "unavailable" - so Hermes can degrade gracefully instead of repeatedly trying to start - a broken local memory backend. - - The embedded daemon computes embeddings via ``sentence_transformers`` - (transformers + huggingface-hub). Importing ``hindsight`` / - ``hindsight_embed`` alone succeeds even when that stack is broken, so - without importing it here the probe would falsely report the backend - healthy and ``hermes memory status`` would stay green while the daemon - aborts at startup on every retain/recall. Import it too so the probe (and - status) reports the real ImportError. - """ - try: - importlib.import_module("hindsight") - importlib.import_module("hindsight_embed.daemon_embed_manager") - importlib.import_module("sentence_transformers") - return True, None - except Exception as exc: - return False, str(exc) - - -def _local_runtime_hint(reason: str | None) -> str: - """Actionable install guidance when the local_embedded runtime is missing. - - ``local_embedded`` imports ``from hindsight import HindsightEmbedded``, which - is provided only by the ``hindsight-all`` package (its wheel ships the - top-level ``hindsight`` module). ``plugin.yaml`` declares only - ``hindsight-client`` (enough for cloud / local_external), so a user who - selected local_embedded without going through ``hermes memory setup`` — a - hand-written config, the legacy ``"mode": "local"`` alias, or a restored - backup — hits ``ModuleNotFoundError: No module named 'hindsight'``. - NousResearch/hermes-agent#7718. - """ - text = (reason or "").lower() - if "no module named" in text and ("hindsight'" in text or 'hindsight"' in text - or "hindsight_embed" in text): - return ( - f" Install the embedded runtime with: uv pip install --python " - f"{sys.executable} hindsight-all — or run 'hermes memory setup'. " - "(local_embedded needs the 'hindsight-all' package, which provides the " - "top-level 'hindsight' module; 'hindsight-client' alone only covers " - "cloud / local_external.)" - ) - return "" - - -def _ensure_cloud_client_dependency() -> None: - """Install the Hindsight cloud client lazily before importing it.""" +def _ensure_client_dependency() -> None: + """Lazily install the Hindsight client (``tools.lazy_deps``) before importing it.""" try: from tools.lazy_deps import ensure as _lazy_ensure _lazy_ensure("memory.hindsight", prompt=False) @@ -213,19 +106,46 @@ def _ensure_cloud_client_dependency() -> None: raise ImportError(str(exc)) from exc +def _cloud_api_key(config: dict) -> str: + return config.get("apiKey") or config.get("api_key") or get_secret("HINDSIGHT_API_KEY", "") + + +def _maybe_upgrade_client() -> None: + """Auto-upgrade an outdated hindsight-client through the environment-aware + lazy_deps installer (sealed hosted venvs redirect to the durable target).""" + try: + from importlib.metadata import version as pkg_version + from packaging.version import Version + installed = pkg_version("hindsight-client") + if Version(installed) < Version(_MIN_CLIENT_VERSION): + logger.warning("hindsight-client %s is outdated (need >=%s), attempting upgrade...", + installed, _MIN_CLIENT_VERSION) + from tools.lazy_deps import install_specs + outcome = install_specs([f"hindsight-client>={_MIN_CLIENT_VERSION}"], timeout=120) + if outcome.ok: + logger.info("hindsight-client upgraded to >=%s", _MIN_CLIENT_VERSION) + elif outcome.blocked: + logger.warning("Auto-upgrade unavailable: %s. Run: uv pip install 'hindsight-client>=%s'", + outcome.reason, _MIN_CLIENT_VERSION) + else: + logger.warning("Auto-upgrade failed: %s. Run: uv pip install 'hindsight-client>=%s'", + (outcome.stderr or "").strip() or "install error", _MIN_CLIENT_VERSION) + except Exception: + pass # packaging not available or other issue — proceed anyway + + # --------------------------------------------------------------------------- -# Hindsight API capability probe — mirrors hindsight-integrations/openclaw. +# Hindsight API capability probe (update_mode='append', Hindsight >= 0.5.0). +# Cached per API URL per process so every provider on the same API shares one +# /version round trip. # --------------------------------------------------------------------------- -# Cache of API_URL -> bool (whether that API supports update_mode='append'). -# Probed once per URL per process — every provider talking to the same API -# gets the same answer without re-hitting /version on each initialize(). _append_capability_cache: Dict[str, bool] = {} _append_capability_lock = threading.Lock() def _meets_minimum_version(actual: str | None, required: str) -> bool: - """Return True if *actual* ≥ *required* (semver). False on missing/invalid.""" + """True if *actual* >= *required* (semver). False on missing/invalid.""" if not actual: return False try: @@ -237,13 +157,8 @@ def _meets_minimum_version(actual: str | None, required: str) -> bool: def _fetch_hindsight_api_version(api_url: str, api_key: str | None = None, timeout: float = 5.0) -> str | None: - """GET ``/version`` and return the version string (or None on failure). - - Hindsight's `/version` endpoint returns ``{"version": "0.5.6", ...}``. - Any failure (timeout, 404, malformed JSON, missing key) → None, which - the caller treats as "legacy API, no update_mode support". - """ - import urllib.error + """GET ``/version`` -> version string, or None on any failure + (the caller treats None as "legacy API, no update_mode support").""" import urllib.request if not api_url: return None @@ -253,8 +168,7 @@ def _fetch_hindsight_api_version(api_url: str, api_key: str | None = None, req.add_header("Authorization", f"Bearer {api_key}") try: with urllib.request.urlopen(req, timeout=timeout) as resp: # noqa: S310 - payload = resp.read().decode("utf-8", errors="replace") - data = json.loads(payload) + data = json.loads(resp.read().decode("utf-8", errors="replace")) except Exception as exc: logger.debug("Hindsight /version probe failed for %s: %s", url, exc) return None @@ -268,9 +182,8 @@ def _check_api_supports_update_mode_append(api_url: str, api_key: str | None = None) -> bool: """Cached capability check for ``update_mode='append'`` on *api_url*. - Probes once per URL per process. Returns False on any probe failure — - that's the safe default: a per-process unique ``document_id`` and no - ``update_mode`` keeps the resume-overwrite fix (#6654) intact. + False on any probe failure — the safe default: a per-process unique + ``document_id`` and no ``update_mode`` keeps the resume-overwrite fix intact. """ if not api_url: return False @@ -280,12 +193,8 @@ def _check_api_supports_update_mode_append(api_url: str, version = _fetch_hindsight_api_version(api_url, api_key) supported = _meets_minimum_version(version, _MIN_VERSION_FOR_UPDATE_MODE_APPEND) with _append_capability_lock: - # Re-check after acquiring the lock in case a concurrent probe filled it. - cached = _append_capability_cache.get(api_url) - if cached is None: - _append_capability_cache[api_url] = supported - else: - supported = cached + # A concurrent probe may have filled the cache meanwhile; its answer wins. + supported = _append_capability_cache.setdefault(api_url, supported) if not supported: logger.warning( "Hindsight API at %s reports version %r, older than %s. " @@ -311,8 +220,7 @@ _loop: asyncio.AbstractEventLoop | None = None _loop_thread: threading.Thread | None = None _loop_lock = threading.Lock() -# Sentinel pushed to the per-provider retain queue to wake the writer for a -# clean exit. A unique object so it can never collide with a real job. +# Pushed to the per-provider retain queue to wake the writer for a clean exit. _WRITER_SENTINEL = object() @@ -336,16 +244,24 @@ def _get_loop() -> asyncio.AbstractEventLoop: def _run_sync(coro, timeout: float = _DEFAULT_TIMEOUT): """Schedule *coro* on the shared loop and block until done.""" from agent.async_utils import safe_schedule_threadsafe - loop = _get_loop() - future = safe_schedule_threadsafe(coro, loop) + future = safe_schedule_threadsafe(coro, _get_loop()) if future is None: raise RuntimeError("Hindsight loop unavailable") return future.result(timeout=timeout) -# --------------------------------------------------------------------------- -# Backward-compatible alias — instances use self._run_sync() instead. -# --------------------------------------------------------------------------- +def _context_thread(target, name: str) -> threading.Thread: + """Daemon thread running *target* in a snapshot of the spawner's contextvars. + + Background threads start with an EMPTY Context; under multiplex_profiles the + spawning thread carries the profile's secret scope + HERMES_HOME override, and + get_secret fails closed without it. (The shared ``hindsight-loop`` needs no + wrap: coroutines scheduled via run_coroutine_threadsafe inherit the + submitter's context per call.) + """ + return threading.Thread( + target=contextvars.copy_context().run, args=(target,), daemon=True, name=name, + ) # --------------------------------------------------------------------------- @@ -419,31 +335,15 @@ REFLECT_SCHEMA = { # --------------------------------------------------------------------------- def _load_config() -> dict: - """Load config from profile-scoped path, legacy path, or env vars. - - Resolution order: - 1. $HERMES_HOME/hindsight/config.json (profile-scoped) - 2. ~/.hindsight/config.json (legacy, shared) - 3. Environment variables - """ - from pathlib import Path - - # Profile-scoped path (preferred) - profile_path = get_hermes_home() / "hindsight" / "config.json" - if profile_path.exists(): - try: - return json.loads(profile_path.read_text(encoding="utf-8")) - except Exception: - pass - - # Legacy shared path (backward compat) - legacy_path = Path.home() / ".hindsight" / "config.json" - if legacy_path.exists(): - try: - return json.loads(legacy_path.read_text(encoding="utf-8")) - except Exception: - pass - + """Load config: $HERMES_HOME/hindsight/config.json (profile-scoped), then + ~/.hindsight/config.json (legacy, shared), then environment variables.""" + for path in (get_hermes_home() / "hindsight" / "config.json", + Path.home() / ".hindsight" / "config.json"): + if path.exists(): + try: + return json.loads(path.read_text(encoding="utf-8")) + except Exception: + pass return { "mode": os.environ.get("HINDSIGHT_MODE", "cloud"), "apiKey": get_secret("HINDSIGHT_API_KEY", ""), @@ -464,300 +364,64 @@ def _load_config() -> dict: } -def _normalize_retain_tags(value: Any) -> List[str]: - """Normalize tag config/tool values to a deduplicated list of strings.""" - if value is None: - return [] - - raw_items: list[Any] - if isinstance(value, list): - raw_items = value - elif isinstance(value, str): - text = value.strip() - if not text: - return [] - if text.startswith("["): - try: - parsed = json.loads(text) - except Exception: - parsed = None - if isinstance(parsed, list): - raw_items = parsed - else: - raw_items = text.split(",") - else: - raw_items = text.split(",") - else: - raw_items = [value] - - normalized = [] - seen = set() - for item in raw_items: - tag = str(item).strip() - if not tag or tag in seen: - continue - seen.add(tag) - normalized.append(tag) - return normalized - - -_OBSERVATION_SCOPE_KEYWORDS = {"per_tag", "combined", "all_combinations"} - - -def _normalize_observation_scopes(value: Any) -> Any: - """Normalize an observation_scopes config value to a Hindsight-accepted form. - - Returns one of: - * ``None`` — nothing configured; Hindsight applies its ``combined`` default. - * a keyword string — ``"per_tag"`` / ``"combined"`` / ``"all_combinations"``. - * ``list[list[str]]`` — custom scopes, one inner list per consolidation pass. - - Accepts a keyword string, a JSON-encoded list, a flat list of tags (treated as - a single scope), or a list of tag-lists. Anything unrecognized yields ``None`` - so we never send an invalid payload. - """ - if value is None: - return None - - if isinstance(value, str): - text = value.strip() - if not text: - return None - if text in _OBSERVATION_SCOPE_KEYWORDS: - return text - if text.startswith("["): - try: - parsed = json.loads(text) - except Exception: - return None - return _normalize_observation_scopes(parsed) - return None - - if isinstance(value, (list, tuple)): - # A flat list of tag strings is one scope; a list of lists is many. - if all(isinstance(entry, str) for entry in value): - inner = [entry.strip() for entry in value if entry.strip()] - return [inner] if inner else None - scopes: list[list[str]] = [] - for entry in value: - if isinstance(entry, (list, tuple)): - inner = [str(tag).strip() for tag in entry if str(tag).strip()] - if inner: - scopes.append(inner) - elif isinstance(entry, str) and entry.strip(): - scopes.append([entry.strip()]) - return scopes or None - - return None - - def _utc_timestamp() -> str: - """Return the UTC write/audit time for retain metadata.""" + """UTC write/audit time for retain metadata.""" return datetime.now(timezone.utc).isoformat(timespec="milliseconds").replace("+00:00", "Z") def _event_timestamp() -> str: - """Return the configured Hermes event time with an explicit UTC offset.""" + """Configured Hermes event time with an explicit UTC offset.""" event_time = _hermes_now() - # hermes_time.now() guarantees an aware datetime. Keep this fallback so a - # replacement clock cannot silently emit an offset-less Hindsight Event Date. + # hermes_time.now() guarantees an aware datetime; the fallback keeps a + # replacement clock from silently emitting an offset-less Event Date. if event_time.tzinfo is None or event_time.utcoffset() is None: event_time = event_time.astimezone() return event_time.isoformat(timespec="seconds") -def _embedded_profile_name(config: dict[str, Any]) -> str: - """Return the Hindsight embedded profile name for this Hermes config.""" - profile = config.get("profile", "hermes") - return str(profile or "hermes") - - -def _load_simple_env(path) -> dict[str, str]: - """Parse a simple KEY=VALUE env file, ignoring comments and blank lines.""" - if not path.exists(): - return {} - - values: dict[str, str] = {} - # utf-8-sig, not plain utf-8: this is also used on the Hermes .env during - # post_setup, and a Notepad BOM would otherwise stick to the first key. - for line in path.read_text(encoding="utf-8-sig", errors="replace").splitlines(): - if not line or line.startswith("#") or "=" not in line: - continue - key, value = line.split("=", 1) - values[key.strip()] = value.strip() - return values - - -def _build_embedded_profile_env(config: dict[str, Any], *, llm_api_key: str | None = None) -> dict[str, str]: - """Build the profile-scoped env file that standalone hindsight-embed consumes.""" - current_key = llm_api_key - if current_key is None: - current_key = ( - config.get("llmApiKey") - or config.get("llm_api_key") - or get_secret("HINDSIGHT_LLM_API_KEY", "") - ) - - current_provider = config.get("llm_provider", "") - current_model = config.get("llm_model", "") - current_base_url = config.get("llm_base_url") or os.environ.get("HINDSIGHT_API_LLM_BASE_URL", "") - - # The embedded daemon expects OpenAI wire format for these providers. - daemon_provider = "openai" if current_provider in {"openai_compatible", "openrouter"} else current_provider - - env_values = { - "HINDSIGHT_API_LLM_PROVIDER": str(daemon_provider), - "HINDSIGHT_API_LLM_API_KEY": str(current_key or ""), - "HINDSIGHT_API_LLM_MODEL": str(current_model), - "HINDSIGHT_API_LOG_LEVEL": "info", - } - if current_base_url: - env_values["HINDSIGHT_API_LLM_BASE_URL"] = str(current_base_url) - - idle_timeout = ( - config.get("idle_timeout") - if config.get("idle_timeout") is not None - else os.environ.get("HINDSIGHT_IDLE_TIMEOUT") - ) - if idle_timeout is not None and idle_timeout != "": - env_values["HINDSIGHT_EMBED_DAEMON_IDLE_TIMEOUT"] = str( - _parse_int_setting(idle_timeout, _DEFAULT_IDLE_TIMEOUT) - ) - return env_values - - -def _embedded_profile_env_path(config: dict[str, Any]): - from pathlib import Path - - return Path.home() / ".hindsight" / "profiles" / f"{_embedded_profile_name(config)}.env" - - -def _secure_write_profile_env(profile_env, content: str) -> None: - """Create/overwrite *profile_env* with owner-only (0600) permissions. - - The file carries the embedded daemon's plaintext LLM API key - (``HINDSIGHT_API_LLM_API_KEY``), so it must never be created with the - default umask-derived mode. A pre-existing file is tightened *before* - the new secret bytes are written. - """ - if profile_env.exists(): - try: - os.chmod(profile_env, 0o600) - except OSError: - pass - fd = os.open(str(profile_env), os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600) - with os.fdopen(fd, "w", encoding="utf-8") as fh: - fh.write(content) - - -def _validate_profile_env_permissions(profile_env) -> None: - """Post-write validation: the secret file must be owner-only on POSIX.""" - if os.name != "posix": - # POSIX mode bits do not model Windows ACLs. - return - import stat - - mode = stat.S_IMODE(profile_env.stat().st_mode) - if mode != 0o600: - try: - os.chmod(profile_env, 0o600) - except OSError: - pass - mode = stat.S_IMODE(profile_env.stat().st_mode) - if mode != 0o600: - raise PermissionError( - f"Embedded Hindsight profile environment is not owner-only: {profile_env}" - ) - - -def _materialize_embedded_profile_env(config: dict[str, Any], *, llm_api_key: str | None = None): - """Write the profile-scoped env file that standalone hindsight-embed uses.""" - profile_env = _embedded_profile_env_path(config) - profile_env.parent.mkdir(parents=True, exist_ok=True) - env_values = _build_embedded_profile_env(config, llm_api_key=llm_api_key) - content = "".join(f"{key}={value}\n" for key, value in env_values.items()) - try: - _secure_write_profile_env(profile_env, content) - _validate_profile_env_permissions(profile_env) - except BaseException: - # Never leave a plaintext API key behind in a file whose permissions - # could not be verified. - try: - profile_env.unlink() - except OSError: - pass - raise - return profile_env - -def _sanitize_bank_segment(value: str) -> str: - """Sanitize a bank_id_template placeholder value. - - Bank IDs should be safe for URL paths and filesystem use. Replaces any - character that isn't alphanumeric, dash, or underscore with a dash, and - collapses runs of dashes. - """ - if not value: - return "" - out = [] - prev_dash = False - for ch in str(value): - if ch.isalnum() or ch == "-" or ch == "_": - out.append(ch) - prev_dash = False - else: - if not prev_dash: - out.append("-") - prev_dash = True - return "".join(out).strip("-_") - - -def _resolve_bank_id_template(template: str, fallback: str, **placeholders: str) -> str: - """Resolve a bank_id template string with the given placeholders. - - Supported placeholders (each is sanitized before substitution): - {profile} — active Hermes profile name (from agent_identity) - {workspace} — Hermes workspace name (from agent_workspace) - {platform} — "cli", "telegram", "discord", etc. - {user} — platform user id (gateway sessions) - {session} — current session id - - Missing/empty placeholders are rendered as the empty string and then - collapsed — e.g. ``hermes-{user}`` with no user becomes ``hermes``. - - If the template is empty, resolution falls back to *fallback*. - Returns the sanitized bank id. - """ - if not template: - return fallback - sanitized = {k: _sanitize_bank_segment(v) for k, v in placeholders.items()} - try: - rendered = template.format(**sanitized) - except (KeyError, IndexError) as exc: - logger.warning("Invalid bank_id_template %r: %s — using fallback %r", - template, exc, fallback) - return fallback - while "--" in rendered: - rendered = rendered.replace("--", "-") - while "__" in rendered: - rendered = rendered.replace("__", "_") - rendered = rendered.strip("-_") - return rendered or fallback +def _mint_document_id(session_id: str) -> str: + """Per-process-lifecycle document id. Reusing session_id alone caused + overwrites on /resume (the reloaded session starts with empty + _session_turns, so its next retain replaced the stored content).""" + return f"{session_id}-{datetime.now().strftime('%Y%m%d_%H%M%S_%f')}" # --------------------------------------------------------------------------- # MemoryProvider implementation # --------------------------------------------------------------------------- +# initialize() kwargs copied verbatim (str, stripped) onto ``self._``. +_SESSION_KWARGS = ( + "platform", "user_id", "user_name", "chat_id", "chat_name", "chat_type", + "thread_id", "agent_identity", "agent_workspace", +) +# Retain metadata keys, each stamped from the attribute of the same name when set. +_METADATA_ATTRS = ( + "session_id", "platform", "user_id", "user_name", "chat_id", "chat_name", + "chat_type", "thread_id", "agent_identity", +) +_SYSTEM_PROMPT_TAILS = { + "context": "Relevant memories are automatically injected into context.", + "tools": ("Use hindsight_recall to search, hindsight_reflect for synthesis, " + "hindsight_retain to store facts."), + "hybrid": ("Relevant memories are automatically injected into context. " + "Use hindsight_recall to search, hindsight_reflect for synthesis, " + "hindsight_retain to store facts."), +} +_TOOL_ERRORS = { + "hindsight_retain": "Failed to store memory", + "hindsight_recall": "Failed to search memory", + "hindsight_reflect": "Failed to reflect", +} + + class HindsightMemoryProvider(MemoryProvider): """Hindsight long-term memory with knowledge graph and multi-strategy retrieval.""" def backup_paths(self) -> List[str]: - """Hindsight's legacy shared config and embedded-mode profile env - files live under ~/.hindsight (see _load_config / line ~509).""" + """Legacy shared config + embedded-mode profile env files live under ~/.hindsight.""" try: - from pathlib import Path - legacy_dir = Path.home() / ".hindsight" - return [str(legacy_dir)] + return [str(Path.home() / ".hindsight")] except Exception: return [] @@ -775,112 +439,77 @@ class HindsightMemoryProvider(MemoryProvider): self._retain_source = _DEFAULT_RETAIN_SOURCE self._retain_user_prefix = "User" self._retain_assistant_prefix = "Assistant" - self._platform = "" - self._user_id = "" - self._user_name = "" - self._chat_id = "" - self._chat_name = "" - self._chat_type = "" - self._thread_id = "" - self._agent_identity = "" - self._agent_workspace = "" + for name in _SESSION_KWARGS: + setattr(self, f"_{name}", "") self._turn_index = 0 self._client = None self._timeout = _DEFAULT_TIMEOUT self._idle_timeout = _DEFAULT_IDLE_TIMEOUT + # Pending prefetch block + its memory count (for the recall indicator). self._prefetch_result = "" - # Number of memories in the pending prefetch block, captured alongside - # _prefetch_result so the deterministic recall indicator can report an - # accurate count without re-parsing the formatted text. self._prefetch_count = 0 self._prefetch_lock = threading.Lock() self._prefetch_thread = None - # State for the model-independent recall indicator (see recall_status()). - # _last_recall_returned tracks whether the most recent prefetch() handed - # any memory to the agent this turn; _last_recall_count is how many. + # Model-independent recall indicator state (see recall_status()). self._last_recall_returned = False self._last_recall_count = 0 self._recall_indicator = True - # Deterministic retain indicator: emitted from sync_turn the moment a - # retain is dispatched to the writer (see _emit_saving_indicator). Uses - # the agent's status channel, injected via initialize(status_callback=). + # Deterministic retain indicator, emitted via the agent's status channel + # (injected through initialize(status_callback=)). self._retain_indicator = True self._status_callback: Optional[Callable[[str], None]] = None - # Single-writer model for retain. sync_turn() enqueues; the writer - # thread drains sequentially. Avoids spawning ad-hoc threads that - # can race the interpreter shutdown and emit "cannot schedule new - # futures after interpreter shutdown" / "Unclosed client session". + # Single-writer model for retain: sync_turn() enqueues, one writer + # thread drains sequentially. Ad-hoc threads raced interpreter shutdown + # ("cannot schedule new futures" / "Unclosed client session"). self._retain_queue: queue.Queue = queue.Queue() self._writer_thread: threading.Thread | None = None self._shutting_down = threading.Event() self._atexit_registered = False - # Server-side async retain operations still in flight. With - # retain_async=True, aretain_batch returns as soon as the write is - # *accepted*, not when it's durable/recall-visible, so the returned - # operation_id(s) stay "pending" until the server finishes. The - # background prefetch gates on these via get_operation_status so recall - # observes the just-completed turn (draining the local queue alone is - # not a read-after-write signal for async retains). + # Server-side async retain ops still in flight. With retain_async=True, + # aretain_batch returns on *acceptance*, not durability, so the + # background prefetch gates on these via get_operation_status (draining + # the local queue alone is not a read-after-write signal). self._pending_retain_ops: set[str] = set() self._pending_retain_ops_lock = threading.Lock() self._retain_ops_bank_id = "" - # Seconds between get_operation_status polls while waiting for server- - # side retain completion. Each poll is a server round trip, so this is - # deliberately coarser than the 0.05s local queue-drain poll: ~20 calls - # max over the default 10s budget instead of ~200. + # Seconds between get_operation_status polls — each is a server round + # trip, so deliberately coarser than the 0.05s local queue-drain poll. self._RETAIN_OP_POLL_INTERVAL_S = 0.5 - # Legacy alias — older tests/callers reference _sync_thread directly. - # Points at _writer_thread once the writer is running. + # Legacy alias — external callers may join _sync_thread; points at the writer. self._sync_thread = None self._session_id = "" self._parent_session_id = "" self._document_id = "" - # Tags self._tags: list[str] | None = None self._recall_tags: list[str] | None = None self._recall_tags_match = "any" - # Retain controls self._auto_retain = True self._retain_every_n_turns = 1 self._retain_async = True - # Async retain never blocks the reply (writes drain on the single - # writer thread). But the next turn's warm prefetch runs on its own - # thread and could read BEFORE the just-completed retain is - # recall-visible on the server, dropping the latest turn from recall. - # When True, the background prefetch first waits (bounded) for the - # local writer queue to drain AND for the server-side async retain - # operation(s) to report completion, an explicit read-after-write - # signal — closing that race without putting any write on the reply - # path. + # Async retain never blocks the reply, but the next turn's warm prefetch + # could read BEFORE the retain is recall-visible. When True the prefetch + # first waits (bounded) for the writer queue to drain AND the server-side + # op(s) to complete — closing the race off the reply path. self._prefetch_waits_for_retain = True self._prefetch_retain_drain_timeout = 10.0 - self._retain_context = "conversation between Hermes Agent and the User" + self._retain_context = _RETAIN_CONTEXT_DEFAULT self._turn_counter = 0 - self._session_turns: list[str] = [] # accumulates ALL turns for the session - # How many turns the last append-mode retain already shipped. Used to - # send only the new delta on subsequent retains when the API supports - # update_mode='append' (legacy/overwrite path still sends everything). + self._session_turns: list[str] = [] # ALL turns for the session + # Turns already shipped by the last append-mode retain (delta watermark). self._last_retained_turn_count = 0 - # Recall controls self._auto_recall = True self._recall_sync = False self._recall_max_tokens = 4096 - # Default to observation-only recall. Observations are Hindsight's - # consolidated knowledge layer — deduplicated, evidence-grounded - # beliefs built from many raw facts, with proof counts and - # freshness signals (see hindsight.vectorize.io/developer/observations). - # Including raw world/experience facts re-ships the supporting - # evidence that observations already summarize, burning the - # `recall_max_tokens` budget. Users can restore the broader - # recall via the `recall_types` config key. + # Observation-only by default: observations are Hindsight's consolidated, + # deduplicated knowledge layer; raw world/experience facts re-ship the + # evidence they summarize and burn the recall_max_tokens budget. self._recall_types: list[str] = ["observation"] self._recall_prompt_preamble = "" self._recall_max_input_chars = 800 - # Bank self._bank_mission = "" self._bank_retain_mission: str | None = None self._bank_id_template = "" @@ -893,48 +522,32 @@ class HindsightMemoryProvider(MemoryProvider): try: cfg = _load_config() mode = cfg.get("mode", "cloud") - if mode in {"local", "local_embedded"}: - available, _ = _check_local_runtime() - return available + if mode in _LOCAL_MODES: + return _check_local_runtime()[0] if mode == "local_external": return True - has_key = bool( - cfg.get("apiKey") - or cfg.get("api_key") - or get_secret("HINDSIGHT_API_KEY", "") - ) - has_url = bool(cfg.get("api_url") or os.environ.get("HINDSIGHT_API_URL", "")) - return has_key or has_url + return bool(_cloud_api_key(cfg)) or bool(cfg.get("api_url") or os.environ.get("HINDSIGHT_API_URL", "")) except Exception: return False def unavailable_reason(self) -> str: - """Explain an unavailable local_embedded provider (missing runtime). - - ``is_available()`` returns False for local modes when the embedded - runtime can't be imported, so ``initialize()`` — and the hint it would - log — is never reached (#7718). Surface the install guidance here, where - agent_init warns about an unavailable provider. - """ + """Install guidance for an unavailable local_embedded runtime. is_available() + gates initialize() out, so the hint it would log is never reached; agent_init + surfaces this instead.""" try: - cfg = _load_config() - mode = cfg.get("mode", "cloud") + mode = _load_config().get("mode", "cloud") except Exception: return "" - if mode not in {"local", "local_embedded"}: + if mode not in _LOCAL_MODES: return "" available, reason = _check_local_runtime() - if available: - return "" - return _local_runtime_hint(reason).strip() + return "" if available else _local_runtime_hint(reason).strip() def save_config(self, values, hermes_home): - """Write config to $HERMES_HOME/hindsight/config.json.""" - import json - from pathlib import Path - config_dir = Path(hermes_home) / "hindsight" - config_dir.mkdir(parents=True, exist_ok=True) - config_path = config_dir / "config.json" + """Merge *values* into $HERMES_HOME/hindsight/config.json.""" + from utils import atomic_json_write + config_path = Path(hermes_home) / "hindsight" / "config.json" + config_path.parent.mkdir(parents=True, exist_ok=True) existing = {} if config_path.exists(): try: @@ -942,247 +555,12 @@ class HindsightMemoryProvider(MemoryProvider): except Exception: pass existing.update(values) - from utils import atomic_json_write atomic_json_write(config_path, existing, mode=0o600) def post_setup(self, hermes_home: str, config: dict) -> None: """Custom setup wizard — installs only the deps needed for the selected mode.""" - import subprocess - import shutil - import sys - from pathlib import Path - - from hermes_cli.config import save_config - from hermes_cli.secret_prompt import masked_secret_prompt - - from hermes_cli.memory_setup import _CANCELLED, _curses_select, _print_cancelled_setup - - print("\n Configuring Hindsight memory:\n") - - existing_config = self._config if isinstance(self._config, dict) else _load_config() - if not isinstance(existing_config, dict): - existing_config = {} - - # Step 1: Mode selection - mode_values = ["cloud", "local_embedded", "local_external"] - mode_items = [ - ("Cloud", "Hindsight Cloud API (lightweight, just needs an API key)"), - ("Local Embedded", "Run Hindsight locally (downloads ~200MB, needs LLM key)"), - ("Local External", "Connect to an existing Hindsight instance"), - ] - existing_mode = existing_config.get("mode") - mode_default_idx = mode_values.index(existing_mode) if existing_mode in mode_values else 0 - mode_idx = _curses_select(" Select mode", mode_items, default=mode_default_idx, cancel_returns=_CANCELLED) - if mode_idx == _CANCELLED: - _print_cancelled_setup() - return - mode = mode_values[mode_idx] - - provider_config: dict = dict(existing_config) - provider_config["mode"] = mode - env_writes: dict = {} - - # Step 2: Install/upgrade deps for selected mode - cloud_dep = f"hindsight-client>={_MIN_CLIENT_VERSION}" - local_dep = "hindsight-all" - if mode == "local_embedded": - deps_to_install = [local_dep] - elif mode == "local_external": - deps_to_install = [cloud_dep] - else: - deps_to_install = [cloud_dep] - - llm_provider = "" - if mode == "local_embedded": - providers_list = list(_PROVIDER_DEFAULT_MODELS.keys()) - llm_items = [ - (p, f"default model: {_PROVIDER_DEFAULT_MODELS[p]}") - for p in providers_list - ] - existing_llm_provider = provider_config.get("llm_provider") - llm_default_idx = providers_list.index(existing_llm_provider) if existing_llm_provider in providers_list else 0 - llm_idx = _curses_select( - " Select LLM provider", - llm_items, - default=llm_default_idx, - cancel_returns=_CANCELLED, - ) - if llm_idx == _CANCELLED: - _print_cancelled_setup() - return - llm_provider = providers_list[llm_idx] - provider_config["llm_provider"] = llm_provider - - print("\n Checking dependencies...") - # Environment-aware install: sealed hosted venvs redirect to the durable - # data-volume target instead of writing to /opt/hermes (NS-605). - from tools.lazy_deps import install_specs - - outcome = install_specs(deps_to_install, timeout=120) - if outcome.ok: - print(" ✓ Dependencies up to date") - elif outcome.blocked: - print(f" ⚠ Cannot install dependencies: {outcome.reason}") - else: - print(f" ⚠ Install failed:\n{(outcome.stderr or '').strip()}") - print(f" Run manually: uv pip install --python {sys.executable} {' '.join(deps_to_install)}") - - # Step 3: Mode-specific config - if mode == "cloud": - print("\n Get your API key at https://ui.hindsight.vectorize.io\n") - existing_key = get_secret("HINDSIGHT_API_KEY", "") or "" - if existing_key: - masked = f"...{existing_key[-4:]}" if len(existing_key) > 4 else "set" - sys.stdout.write(f" API key (current: {masked}, blank to keep): ") - sys.stdout.flush() - api_key = masked_secret_prompt("") if sys.stdin.isatty() else sys.stdin.readline().strip() - else: - sys.stdout.write(" API key: ") - sys.stdout.flush() - api_key = masked_secret_prompt("") if sys.stdin.isatty() else sys.stdin.readline().strip() - if api_key: - env_writes["HINDSIGHT_API_KEY"] = api_key - - val = input(f" API URL [{_DEFAULT_API_URL}]: ").strip() - if val: - provider_config["api_url"] = val - - elif mode == "local_external": - val = input(f" Hindsight API URL [{_DEFAULT_LOCAL_URL}]: ").strip() - provider_config["api_url"] = val or _DEFAULT_LOCAL_URL - - sys.stdout.write(" API key (optional, blank to skip): ") - sys.stdout.flush() - api_key = masked_secret_prompt("") if sys.stdin.isatty() else sys.stdin.readline().strip() - if api_key: - env_writes["HINDSIGHT_API_KEY"] = api_key - - else: # local_embedded - if llm_provider == "openai_compatible": - existing_base_url = provider_config.get("llm_base_url", "") - prompt = " LLM endpoint URL (e.g. http://192.168.1.10:8080/v1)" - if existing_base_url: - prompt += f" [{existing_base_url}]" - prompt += ": " - val = input(prompt).strip() - if val: - provider_config["llm_base_url"] = val - elif llm_provider == "openrouter": - provider_config["llm_base_url"] = "https://openrouter.ai/api/v1" - - provider_default_model = _PROVIDER_DEFAULT_MODELS.get(llm_provider, "gpt-4o-mini") - current_model = provider_config.get("llm_model") or provider_default_model - val = input(f" LLM model [{current_model}]: ").strip() - provider_config["llm_model"] = val or current_model - - sys.stdout.write(" LLM API key: ") - sys.stdout.flush() - llm_key = masked_secret_prompt("") if sys.stdin.isatty() else sys.stdin.readline().strip() - if llm_key: - env_writes["HINDSIGHT_LLM_API_KEY"] = llm_key - else: - env_path = Path(hermes_home) / ".env" - existing_llm_key = "" - if env_path.exists(): - # utf-8-sig: a Notepad BOM must not hide the first key. - for line in env_path.read_text(encoding="utf-8-sig").splitlines(): - if line.startswith("HINDSIGHT_LLM_API_KEY="): - existing_llm_key = line.split("=", 1)[1] - break - env_writes["HINDSIGHT_LLM_API_KEY"] = existing_llm_key - - # Step 4: Save everything - provider_config.setdefault("bank_id", "hermes") - provider_config.setdefault("recall_budget", "mid") - # Read existing timeout from config if present, otherwise use default. - # Preserve explicit 0 values instead of treating them as blank. - existing_timeout = provider_config.get("timeout") - timeout_val = existing_timeout if existing_timeout is not None else _DEFAULT_TIMEOUT - provider_config["timeout"] = timeout_val - env_writes["HINDSIGHT_TIMEOUT"] = str(timeout_val) - if mode == "local_embedded": - existing_idle_timeout = provider_config.get("idle_timeout") - idle_timeout_val = existing_idle_timeout if existing_idle_timeout is not None else _DEFAULT_IDLE_TIMEOUT - provider_config["idle_timeout"] = idle_timeout_val - env_writes["HINDSIGHT_IDLE_TIMEOUT"] = str(idle_timeout_val) - config["memory"]["provider"] = "hindsight" - save_config(config) - - self.save_config(provider_config, hermes_home) - - if env_writes: - env_path = Path(hermes_home) / ".env" - env_path.parent.mkdir(parents=True, exist_ok=True) - existing_lines = [] - if env_path.exists(): - # utf-8-sig: a Notepad BOM would glue U+FEFF onto the first - # key, defeating the in-place update below and appending a - # duplicate line instead. - existing_lines = env_path.read_text(encoding="utf-8-sig").splitlines() - updated_keys = set() - new_lines = [] - for line in existing_lines: - key_match = line.split("=", 1)[0].strip() if "=" in line and not line.startswith("#") else None - if key_match and key_match in env_writes: - new_lines.append(f"{key_match}={env_writes[key_match]}") - updated_keys.add(key_match) - else: - new_lines.append(line) - for k, v in env_writes.items(): - if k not in updated_keys: - new_lines.append(f"{k}={v}") - env_path.write_text("\n".join(new_lines) + "\n", encoding="utf-8") - - # Step 5: Optional starter template. Only for cloud / local_external — - # the API is reachable now; local_embedded's daemon isn't up during setup. - from . import templates as _hs_templates - if _hs_templates.supported_for_mode(mode): - self._offer_starter_template(mode, provider_config, env_writes) - - if mode == "local_embedded": - materialized_config = dict(provider_config) - config_path = Path(hermes_home) / "hindsight" / "config.json" - try: - materialized_config = json.loads(config_path.read_text(encoding="utf-8")) - except Exception: - pass - - llm_api_key = env_writes.get("HINDSIGHT_LLM_API_KEY", "") - if not llm_api_key: - llm_api_key = _load_simple_env(Path(hermes_home) / ".env").get("HINDSIGHT_LLM_API_KEY", "") - if not llm_api_key: - llm_api_key = _load_simple_env(_embedded_profile_env_path(materialized_config)).get( - "HINDSIGHT_API_LLM_API_KEY", - "", - ) - - _materialize_embedded_profile_env( - materialized_config, - llm_api_key=llm_api_key or None, - ) - - print(f"\n ✓ Hindsight memory configured ({mode} mode)") - if env_writes: - print(" API keys saved to .env") - print("\n Start a new session to activate.\n") - - def _offer_starter_template(self, mode: str, provider_config: dict, env_writes: dict) -> None: - """Offer to seed the bank with a Hermes starter template (best-effort).""" - from hermes_cli.memory_setup import _CANCELLED, _curses_select - - from . import templates as _hs_templates - - default_url = _DEFAULT_LOCAL_URL if mode == "local_external" else _DEFAULT_API_URL - api_url = provider_config.get("api_url") or default_url - bank_id = provider_config.get("bank_id", "hermes") - api_key = env_writes.get("HINDSIGHT_API_KEY") or os.environ.get("HINDSIGHT_API_KEY", "") or None - _hs_templates.run_template_step( - api_url=api_url, - bank_id=bank_id, - api_key=api_key, - select=_curses_select, - cancelled=_CANCELLED, - ) + from .setup import run_setup + run_setup(self, hermes_home, config) def get_config_schema(self): return [ @@ -1231,279 +609,105 @@ class HindsightMemoryProvider(MemoryProvider): {"key": "port_health_grace_timeout", "description": "Seconds to wait for a slow daemon /health before treating it as stale (raise on busy/low-resource hosts; blank uses the 30s default)", "default": "", "when": {"mode": "local_embedded"}}, ] + # -- client ------------------------------------------------------------- + + def _int_setting(self, key: str, env_var: str, default: int, env_default=None) -> int: + """Config value if set (explicit 0 preserved), else env var, else default.""" + value = self._config.get(key) + if value is None: + value = os.environ.get(env_var, env_default) + return _parse_int_setting(value, default) + + def _new_embedded_client(self): + available, reason = _check_local_runtime() + if not available: + raise RuntimeError("Hindsight local runtime is unavailable" + (f": {reason}" if reason else "")) + _ensure_client_dependency() + from hindsight import HindsightEmbedded + HindsightEmbedded.__del__ = lambda self: None + cfg = self._config + llm_provider = _daemon_llm_provider(cfg.get("llm_provider", "")) + logger.debug("Creating HindsightEmbedded client (profile=%s, provider=%s)", + cfg.get("profile", "hermes"), llm_provider) + kwargs = dict( + profile=cfg.get("profile", "hermes"), + llm_provider=llm_provider, + llm_api_key=cfg.get("llmApiKey") or cfg.get("llm_api_key") or get_secret("HINDSIGHT_LLM_API_KEY", ""), + llm_model=cfg.get("llm_model", ""), + ) + if self._llm_base_url: + kwargs["llm_base_url"] = self._llm_base_url + self._idle_timeout = self._int_setting( + "idle_timeout", "HINDSIGHT_IDLE_TIMEOUT", _DEFAULT_IDLE_TIMEOUT, env_default=self._idle_timeout, + ) + kwargs["idle_timeout"] = self._idle_timeout + return HindsightEmbedded(**kwargs) + + def _new_cloud_client(self): + _ensure_client_dependency() + from hindsight_client import Hindsight + kwargs = {"base_url": self._api_url, "timeout": float(self._timeout or _DEFAULT_TIMEOUT)} + if self._api_key: + kwargs["api_key"] = self._api_key + logger.debug("Creating Hindsight cloud client (url=%s, has_key=%s, timeout=%s)", + self._api_url, bool(self._api_key), kwargs["timeout"]) + return Hindsight(**kwargs) + def _get_client(self): """Return the cached Hindsight client (created once, reused).""" if self._client is None: - if self._mode == "local_embedded": - available, reason = _check_local_runtime() - if not available: - raise RuntimeError( - "Hindsight local runtime is unavailable" - + (f": {reason}" if reason else "") - ) - try: - from tools.lazy_deps import ensure as _lazy_ensure - _lazy_ensure("memory.hindsight", prompt=False) - except ImportError: - pass - except Exception as _e: - raise ImportError(str(_e)) - from hindsight import HindsightEmbedded - HindsightEmbedded.__del__ = lambda self: None - llm_provider = self._config.get("llm_provider", "") - if llm_provider in {"openai_compatible", "openrouter"}: - llm_provider = "openai" - logger.debug("Creating HindsightEmbedded client (profile=%s, provider=%s)", - self._config.get("profile", "hermes"), llm_provider) - kwargs = dict( - profile=self._config.get("profile", "hermes"), - llm_provider=llm_provider, - llm_api_key=self._config.get("llmApiKey") or self._config.get("llm_api_key") or get_secret("HINDSIGHT_LLM_API_KEY", ""), - llm_model=self._config.get("llm_model", ""), - ) - if self._llm_base_url: - kwargs["llm_base_url"] = self._llm_base_url - idle_timeout = _parse_int_setting( - self._config.get("idle_timeout") - if self._config.get("idle_timeout") is not None - else os.environ.get("HINDSIGHT_IDLE_TIMEOUT", self._idle_timeout), - _DEFAULT_IDLE_TIMEOUT, - ) - self._idle_timeout = idle_timeout - kwargs["idle_timeout"] = idle_timeout - self._client = HindsightEmbedded(**kwargs) - else: - _ensure_cloud_client_dependency() - from hindsight_client import Hindsight - timeout = self._timeout or _DEFAULT_TIMEOUT - kwargs = {"base_url": self._api_url, "timeout": float(timeout)} - if self._api_key: - kwargs["api_key"] = self._api_key - logger.debug("Creating Hindsight cloud client (url=%s, has_key=%s, timeout=%s)", - self._api_url, bool(self._api_key), kwargs["timeout"]) - self._client = Hindsight(**kwargs) + self._client = ( + self._new_embedded_client() if self._mode == "local_embedded" else self._new_cloud_client() + ) return self._client def _run_sync(self, coro): """Schedule *coro* on the shared loop using the configured timeout.""" return _run_sync(coro, timeout=self._timeout) - def _is_retriable_embedded_connection_error(self, exc: Exception) -> bool: - """Return True for stale embedded-daemon connection failures.""" - if self._mode != "local_embedded": - return False - text = f"{type(exc).__name__}: {exc}".lower() - return any( - marker in text - for marker in ( - "cannot connect to host", - "connection refused", - "connect call failed", - "clientconnectorerror", + def _run_hindsight_operation(self, operation): + """Run an async client operation; for local_embedded, a stale-daemon + connection failure recreates the client and retries once.""" + client = self._get_client() + try: + return self._run_sync(operation(client)) + except Exception as exc: + text = f"{type(exc).__name__}: {exc}".lower() + if self._mode != "local_embedded" or not any(m in text for m in _RETRIABLE_CONNECTION_MARKERS): + raise + logger.info( + "Hindsight embedded daemon appears unreachable; recreating client and retrying once: %s", + exc, ) - ) + self._client = None + client = self._get_client() + self._client = client + return self._run_sync(operation(client)) + + # -- retain writer thread + server-side visibility ------------------------- def _ensure_writer(self) -> None: - """Lazy-start the single retain-writer thread. - - We don't start the writer in initialize() so providers that never - retain (e.g. tools-only mode) don't pay for an idle thread. - """ + """Lazy-start the single retain-writer thread (providers that never + retain, e.g. tools-only mode, don't pay for an idle thread).""" thread = self._writer_thread if thread is not None and thread.is_alive(): return - # If the previous writer exited (e.g. after a prior shutdown), reset - # the flag so this fresh writer is allowed to drain new jobs. + # A previous writer may have exited after shutdown(); allow the fresh one to drain. self._shutting_down.clear() - # Per-provider background threads start with an EMPTY contextvars - # Context. Under multiplex_profiles the spawning thread carries the - # profile's secret scope + HERMES_HOME override (gateway/run.py wraps - # the agent turn in copy_context().run), and get_secret fails closed - # without it (#92608). Snapshot the spawner's context into the thread. - # (The shared ``hindsight-loop`` thread needs no wrap: coroutines - # scheduled via run_coroutine_threadsafe inherit the submitter's - # context per call, so one loop can serve every profile.) - thread = threading.Thread( - target=contextvars.copy_context().run, - args=(self._writer_loop,), - daemon=True, - name="hindsight-writer", - ) - self._writer_thread = thread - # Keep the legacy _sync_thread alias pointing at the writer so any - # external code that joins _sync_thread keeps working. - self._sync_thread = thread + thread = _context_thread(self._writer_loop, "hindsight-writer") + self._writer_thread = self._sync_thread = thread thread.start() - def _track_retain_ops(self, retain_response, bank_id: str) -> None: - """Record server-side async operation id(s) from an aretain_batch reply. - - Async retains return ``operation_id`` / ``operation_ids`` that stay - ``pending`` on the server until the write is durable and recall-visible. - The bank_id is captured alongside so completion can be polled with the - same bank the write targeted. - """ - ids: list[str] = [] - single = getattr(retain_response, "operation_id", None) - if single: - ids.append(str(single)) - multiple = getattr(retain_response, "operation_ids", None) - if multiple: - ids.extend(str(op) for op in multiple if op) - if not ids: - # Server didn't hand back an op id (older API, or it completed - # synchronously). Nothing to poll — local queue drain is the only - # available signal in that case. - return - self._retain_ops_bank_id = bank_id - with self._pending_retain_ops_lock: - self._pending_retain_ops.update(ids) - - def _is_retain_op_complete(self, bank_id: str, op_id: str) -> bool: - """Return True when a server-side async retain op is done (or gone). - - ``get_operation_status`` returns ``completed``/``failed`` for a known - op; completed ops are evicted server-side, so a NotFound (404) also - means "no longer pending" and is treated as done. Transient errors - return False so the caller keeps waiting until its deadline. - """ - from hindsight_client_api.exceptions import NotFoundException - - try: - resp = self._run_hindsight_operation( - lambda client: client.operations.get_operation_status( - bank_id=bank_id, operation_id=op_id - ) - ) - except NotFoundException: - return True - except Exception as exc: - logger.debug("Prefetch: operation status check failed for %s: %s", op_id, exc) - return False - status = str(getattr(resp, "status", "") or "").lower() - return status in {"completed", "failed"} - - def _wait_for_retains_drained(self, timeout: float) -> bool: - """Block up to *timeout* seconds for the just-completed turn's retain to - become recall-visible on the server. - - Used by the background prefetch so the next turn's recall observes the - just-completed turn's write instead of racing ahead of it. Runs only on - the background prefetch thread — never on the reply path. - - Two ordered barriers, both bounded by the shared *timeout* budget: - - 1. Local writer queue drains (the retain call has been *dispatched* to - the server). Polls ``unfinished_tasks`` rather than ``queue.join()`` - so a wedged write can't hang the prefetch. - 2. Server-side async operations complete. With ``retain_async=True`` the - dispatched call returns on *acceptance*, not durability, so draining - the local queue alone is NOT a read-after-write signal. We poll - ``get_operation_status`` for the tracked op id(s) until the server - reports completion (an explicit read-after-write condition). - - Returns True if both barriers cleared within the budget, False on - timeout/shutdown. - """ - deadline = None if timeout <= 0 else time.monotonic() + timeout - - def _expired() -> bool: - return deadline is not None and time.monotonic() >= deadline - - # Barrier 1: local queue drain (retain dispatched to the server). - while self._retain_queue.unfinished_tasks > 0: - if self._shutting_down.is_set(): - return False - if _expired(): - logger.debug( - "Prefetch: retain drain timed out after %.1fs (%d pending)", - timeout, self._retain_queue.unfinished_tasks, - ) - return False - time.sleep(0.05) - - # Barrier 2: server-side async retain completion (read-after-write). - return self._wait_for_server_retain_ops(deadline, timeout) - - def _wait_for_server_retain_ops(self, deadline: float | None, timeout: float) -> bool: - """Poll tracked async retain ops until complete or the deadline passes. - - *deadline* is a ``time.monotonic()`` value (None = no bound). Completed - ops are removed from the pending set as they finish so a later prefetch - doesn't re-poll them. - - Ops still pending when the deadline expires are DROPPED, not retained: - keeping them would make a permanently failing status endpoint (auth - error, endless 500s, server that loses ops without a 404) grow the - pending set forever and burn the full timeout on EVERY subsequent - prefetch — turning "bounded wait per prefetch" into unbounded - session-wide degradation (and, via prefetch()'s bounded join on the - reply path, a per-turn reply-latency penalty). Dropping trades a - possibly-stale recall NOW (identical to prefetch_waits_for_retain=False - behavior) for guaranteed liveness; the drop is logged at WARNING once - per prefetch so persistent server trouble is visible. - - Status polls are spaced by _RETAIN_OP_POLL_INTERVAL_S (0.5s) — server - round trips per op are bounded (~20 over a 10s budget), unlike the - cheap 0.05s local queue-drain poll in _wait_for_retains_drained. - """ - while True: - with self._pending_retain_ops_lock: - bank_id = getattr(self, "_retain_ops_bank_id", "") or self._bank_id - pending = list(self._pending_retain_ops) - if not pending: - return True - if self._shutting_down.is_set(): - return False - - done: set[str] = set() - expired = False - for op_id in pending: - if self._shutting_down.is_set(): - return False - if deadline is not None and time.monotonic() >= deadline: - expired = True - break - if self._is_retain_op_complete(bank_id, op_id): - done.add(op_id) - - if expired: - with self._pending_retain_ops_lock: - self._pending_retain_ops.difference_update(done) - dropped = len(self._pending_retain_ops) - self._pending_retain_ops.clear() - logger.warning( - "Prefetch: server retain visibility timed out after %.1fs; " - "dropping %d unresolved op(s) so later prefetches stay " - "bounded (recall may miss the just-completed turn)", - timeout, dropped, - ) - return False - - with self._pending_retain_ops_lock: - self._pending_retain_ops.difference_update(done) - still_pending = bool(self._pending_retain_ops) - if not still_pending: - return True - if deadline is not None and time.monotonic() >= deadline: - with self._pending_retain_ops_lock: - dropped = len(self._pending_retain_ops) - self._pending_retain_ops.clear() - logger.warning( - "Prefetch: server retain visibility timed out after %.1fs; " - "dropping %d unresolved op(s) so later prefetches stay " - "bounded (recall may miss the just-completed turn)", - timeout, dropped, - ) - return False - time.sleep(self._RETAIN_OP_POLL_INTERVAL_S) + def _register_atexit(self) -> None: + """Idempotent atexit drain so a CLI exit that skips + MemoryManager.shutdown_all() can't race interpreter teardown.""" + if not self._atexit_registered: + self._atexit_registered = True + atexit.register(self._atexit_shutdown) def _writer_loop(self) -> None: - """Drain the retain queue serially. Exits on sentinel. - - Each job() is wrapped so a single failure can't kill the writer. - task_done() always fires so queue.join() works in tests. - """ + """Drain the retain queue serially; exits on the sentinel. A failing job + can't kill the writer, and task_done() always fires so queue.join() works.""" while True: try: job = self._retain_queue.get(timeout=1.0) @@ -1521,19 +725,6 @@ class HindsightMemoryProvider(MemoryProvider): finally: self._retain_queue.task_done() - def _register_atexit(self) -> None: - """Register an idempotent atexit hook to drain the writer. - - Without this, a CLI exit that doesn't go through MemoryManager. - shutdown_all() would leave in-flight retain jobs racing interpreter - teardown, producing "cannot schedule new futures" warnings and - unclosed aiohttp sessions. - """ - if self._atexit_registered: - return - self._atexit_registered = True - atexit.register(self._atexit_shutdown) - def _atexit_shutdown(self) -> None: if self._shutting_down.is_set(): return @@ -1542,30 +733,112 @@ class HindsightMemoryProvider(MemoryProvider): except Exception as exc: logger.debug("Hindsight atexit shutdown failed: %s", exc) - def _run_hindsight_operation(self, operation): - """Run an async Hindsight client operation, retrying once after idle shutdown.""" - client = self._get_client() + def _track_retain_ops(self, retain_response, bank_id: str) -> None: + """Record the async ``operation_id``/``operation_ids`` from an + aretain_batch reply (pending server-side until recall-visible), with the + bank they were written to. No id (older API / synchronous completion) + means only the local queue drain is available as a signal.""" + ids: list[str] = [] + single = getattr(retain_response, "operation_id", None) + if single: + ids.append(str(single)) + ids.extend(str(op) for op in (getattr(retain_response, "operation_ids", None) or []) if op) + if not ids: + return + self._retain_ops_bank_id = bank_id + with self._pending_retain_ops_lock: + self._pending_retain_ops.update(ids) + + def _is_retain_op_complete(self, bank_id: str, op_id: str) -> bool: + """True when a server-side retain op is done or gone. Completed ops are + evicted server-side, so NotFound (404) also means "no longer pending". + Transient errors return False so the caller keeps waiting.""" + from hindsight_client_api.exceptions import NotFoundException + try: - return self._run_sync(operation(client)) - except Exception as exc: - if not self._is_retriable_embedded_connection_error(exc): - raise - logger.info( - "Hindsight embedded daemon appears unreachable; recreating client and retrying once: %s", - exc, + resp = self._run_hindsight_operation( + lambda client: client.operations.get_operation_status(bank_id=bank_id, operation_id=op_id) ) - self._client = None - client = self._get_client() - self._client = client - return self._run_sync(operation(client)) + except NotFoundException: + return True + except Exception as exc: + logger.debug("Prefetch: operation status check failed for %s: %s", op_id, exc) + return False + return str(getattr(resp, "status", "") or "").lower() in {"completed", "failed"} + + def _wait_for_retains_drained(self, timeout: float) -> bool: + """Block up to *timeout* seconds for the just-completed turn's retain to + become recall-visible. Background prefetch thread only — never the reply path. + + Two ordered barriers on one budget: (1) the local writer queue drains + (retain *dispatched*) — polls ``unfinished_tasks`` rather than + ``queue.join()`` so a wedged write can't hang the prefetch; (2) the + server-side async ops complete (an explicit read-after-write signal, + since async retain returns on acceptance, not durability). + Returns False on timeout/shutdown. + """ + deadline = None if timeout <= 0 else time.monotonic() + timeout + while self._retain_queue.unfinished_tasks > 0: + if self._shutting_down.is_set(): + return False + if deadline is not None and time.monotonic() >= deadline: + logger.debug("Prefetch: retain drain timed out after %.1fs (%d pending)", + timeout, self._retain_queue.unfinished_tasks) + return False + time.sleep(0.05) + return self._wait_for_server_retain_ops(deadline, timeout) + + def _wait_for_server_retain_ops(self, deadline: float | None, timeout: float) -> bool: + """Poll tracked async retain ops until complete or *deadline* (monotonic; + None = unbounded). Completed ops leave the pending set as they finish. + + Ops still pending at the deadline are DROPPED: keeping them would let a + permanently failing status endpoint grow the pending set forever and burn + the full timeout on EVERY later prefetch (and, via prefetch()'s bounded + join, a per-turn reply-latency penalty). Dropping trades a possibly-stale + recall now for guaranteed liveness; logged at WARNING once per prefetch. + """ + def _expired() -> bool: + return deadline is not None and time.monotonic() >= deadline + + while True: + with self._pending_retain_ops_lock: + bank_id = self._retain_ops_bank_id or self._bank_id + pending = list(self._pending_retain_ops) + if not pending: + return True + if self._shutting_down.is_set(): + return False + done: set[str] = set() + for op_id in pending: + if self._shutting_down.is_set(): + return False + if _expired(): + break + if self._is_retain_op_complete(bank_id, op_id): + done.add(op_id) + with self._pending_retain_ops_lock: + self._pending_retain_ops.difference_update(done) + if not self._pending_retain_ops: + return True + dropped = len(self._pending_retain_ops) if _expired() else 0 + if dropped: + self._pending_retain_ops.clear() + if dropped: + logger.warning( + "Prefetch: server retain visibility timed out after %.1fs; " + "dropping %d unresolved op(s) so later prefetches stay " + "bounded (recall may miss the just-completed turn)", + timeout, dropped, + ) + return False + time.sleep(self._RETAIN_OP_POLL_INTERVAL_S) + + # -- retain target ----------------------------------------------------------- def _probe_url(self) -> str: - """Return the URL to probe /version on. - - For local_embedded the daemon is on a per-profile dynamic port, - so we prefer the running client's URL when available; otherwise - fall back to the configured api_url. - """ + """URL to probe /version on: the running embedded client's dynamic + per-profile port when available, else the configured api_url.""" if self._mode == "local_embedded" and self._client is not None: url = getattr(self._client, "url", None) if url: @@ -1573,201 +846,133 @@ class HindsightMemoryProvider(MemoryProvider): return self._api_url or "" def _resolve_retain_target(self, fallback_document_id: str) -> tuple[str, str | None]: - """Pick (document_id, update_mode) based on live API capability. + """Pick (document_id, update_mode) from live API capability. - On Hindsight ≥ 0.5.0 the API supports ``update_mode='append'``, - which lets us reuse a stable session-scoped ``document_id`` across - process lifecycles without overwriting prior turns. On older APIs - we fall back to *fallback_document_id* (the per-process unique - ``f"{session_id}-{start_ts}"`` minted at initialize / switch time) - and don't pass ``update_mode`` at all — that's the only way the - resume-overwrite fix (#6654) keeps working on legacy servers. - - Probe is cached at module level per API URL, so this is one HTTP - round-trip per (process, api_url) pair regardless of how many - retains fire. + Hindsight >= 0.5.0 supports ``update_mode='append'``, so the stable + session-scoped ``document_id`` can be reused across process lifecycles + without overwriting prior turns. Older APIs get *fallback_document_id* + (per-process unique, minted at initialize/switch time) and no + ``update_mode`` — the only way the resume-overwrite fix works there. + The probe is cached per (process, api_url). """ - if not self._session_id: - return fallback_document_id, None - if _check_api_supports_update_mode_append(self._probe_url(), self._api_key): + if self._session_id and _check_api_supports_update_mode_append(self._probe_url(), self._api_key): return self._session_id, "append" return fallback_document_id, None + # -- lifecycle --------------------------------------------------------------- + def initialize(self, session_id: str, **kwargs) -> None: self._session_id = str(session_id or "").strip() self._parent_session_id = str(kwargs.get("parent_session_id", "") or "").strip() - # Agent status channel for the deterministic retain indicator (recall - # emits via the pull-based recall_status()/describe_recall() path). - _status_cb = kwargs.get("status_callback") - if callable(_status_cb): - self._status_callback = _status_cb + # Status channel for the retain indicator (recall reports via recall_status()). + status_cb = kwargs.get("status_callback") + if callable(status_cb): + self._status_callback = status_cb + # session_id stays in tags so processes for one session remain filterable together. + self._document_id = _mint_document_id(self._session_id) + _maybe_upgrade_client() - # Each process lifecycle gets its own document_id. Reusing session_id - # alone caused overwrites on /resume — the reloaded session starts - # with an empty _session_turns, so the next retain would replace the - # previously stored content. session_id stays in tags so processes - # for the same session remain filterable together. - start_ts = datetime.now().strftime("%Y%m%d_%H%M%S_%f") - self._document_id = f"{self._session_id}-{start_ts}" - - # Check client version and auto-upgrade if needed - try: - from importlib.metadata import version as pkg_version - from packaging.version import Version - installed = pkg_version("hindsight-client") - if Version(installed) < Version(_MIN_CLIENT_VERSION): - logger.warning("hindsight-client %s is outdated (need >=%s), attempting upgrade...", - installed, _MIN_CLIENT_VERSION) - # Environment-aware install: sealed hosted venvs redirect to the - # durable data-volume target instead of /opt/hermes (NS-605). - from tools.lazy_deps import install_specs - outcome = install_specs([f"hindsight-client>={_MIN_CLIENT_VERSION}"], timeout=120) - if outcome.ok: - logger.info("hindsight-client upgraded to >=%s", _MIN_CLIENT_VERSION) - elif outcome.blocked: - logger.warning("Auto-upgrade unavailable: %s. Run: uv pip install 'hindsight-client>=%s'", - outcome.reason, _MIN_CLIENT_VERSION) - else: - logger.warning("Auto-upgrade failed: %s. Run: uv pip install 'hindsight-client>=%s'", - (outcome.stderr or "").strip() or "install error", _MIN_CLIENT_VERSION) - except Exception: - pass # packaging not available or other issue — proceed anyway - - self._config = _load_config() - self._platform = str(kwargs.get("platform") or "").strip() - self._user_id = str(kwargs.get("user_id") or "").strip() - self._user_name = str(kwargs.get("user_name") or "").strip() - self._chat_id = str(kwargs.get("chat_id") or "").strip() - self._chat_name = str(kwargs.get("chat_name") or "").strip() - self._chat_type = str(kwargs.get("chat_type") or "").strip() - self._thread_id = str(kwargs.get("thread_id") or "").strip() - self._agent_identity = str(kwargs.get("agent_identity") or "").strip() - self._agent_workspace = str(kwargs.get("agent_workspace") or "").strip() + self._config = cfg = _load_config() + for name in _SESSION_KWARGS: + setattr(self, f"_{name}", str(kwargs.get(name) or "").strip()) self._turn_index = 0 self._session_turns = [] self._last_retained_turn_count = 0 - self._mode = self._config.get("mode", "cloud") - # Read timeout from config or env var, fall back to default - self._timeout = _parse_int_setting( - self._config.get("timeout") if self._config.get("timeout") is not None else os.environ.get("HINDSIGHT_TIMEOUT"), - _DEFAULT_TIMEOUT, - ) - self._idle_timeout = _parse_int_setting( - self._config.get("idle_timeout") if self._config.get("idle_timeout") is not None else os.environ.get("HINDSIGHT_IDLE_TIMEOUT"), - _DEFAULT_IDLE_TIMEOUT, - ) - # "local" is a legacy alias for "local_embedded" - if self._mode == "local": + self._mode = cfg.get("mode", "cloud") + self._timeout = self._int_setting("timeout", "HINDSIGHT_TIMEOUT", _DEFAULT_TIMEOUT) + self._idle_timeout = self._int_setting("idle_timeout", "HINDSIGHT_IDLE_TIMEOUT", _DEFAULT_IDLE_TIMEOUT) + if self._mode == "local": # legacy alias self._mode = "local_embedded" if self._mode == "local_embedded": - # Export the daemon health grace timeout BEFORE importing - # daemon_embed_manager (which reads it at import time). - _export_port_health_grace_timeout(self._config) + # Must precede the daemon_embed_manager import, which reads it at import time. + _export_port_health_grace_timeout(cfg) available, reason = _check_local_runtime() if not available: logger.warning( "Hindsight local mode disabled because its runtime could not be imported: %s.%s", - reason, - _local_runtime_hint(reason), + reason, _local_runtime_hint(reason), ) self._mode = "disabled" return - self._api_key = self._config.get("apiKey") or self._config.get("api_key") or get_secret("HINDSIGHT_API_KEY", "") + self._api_key = _cloud_api_key(cfg) default_url = _DEFAULT_LOCAL_URL if self._mode in {"local_embedded", "local_external"} else _DEFAULT_API_URL - self._api_url = self._config.get("api_url") or os.environ.get("HINDSIGHT_API_URL", default_url) - self._llm_base_url = self._config.get("llm_base_url", "") + self._api_url = cfg.get("api_url") or os.environ.get("HINDSIGHT_API_URL", default_url) + self._llm_base_url = cfg.get("llm_base_url", "") - banks = cfg_get(self._config, "banks", "hermes", default={}) - static_bank_id = self._config.get("bank_id") or banks.get("bankId", "hermes") - self._bank_id_template = self._config.get("bank_id_template", "") or "" + banks = cfg_get(cfg, "banks", "hermes", default={}) + self._bank_id_template = cfg.get("bank_id_template", "") or "" self._bank_id = _resolve_bank_id_template( self._bank_id_template, - fallback=static_bank_id, + fallback=cfg.get("bank_id") or banks.get("bankId", "hermes"), profile=self._agent_identity, workspace=self._agent_workspace, platform=self._platform, user=self._user_id, session=self._session_id, ) - budget = self._config.get("recall_budget") or self._config.get("budget") or banks.get("budget", "mid") + budget = cfg.get("recall_budget") or cfg.get("budget") or banks.get("budget", "mid") self._budget = budget if budget in _VALID_BUDGETS else "mid" - - memory_mode = self._config.get("memory_mode", "hybrid") - self._memory_mode = memory_mode if memory_mode in {"context", "tools", "hybrid"} else "hybrid" - - prefetch_method = self._config.get("recall_prefetch_method") or self._config.get("prefetch_method", "recall") + memory_mode = cfg.get("memory_mode", "hybrid") + self._memory_mode = memory_mode if memory_mode in _SYSTEM_PROMPT_TAILS else "hybrid" + prefetch_method = cfg.get("recall_prefetch_method") or cfg.get("prefetch_method", "recall") self._prefetch_method = prefetch_method if prefetch_method in {"recall", "reflect"} else "recall" - # Bank options - self._bank_mission = self._config.get("bank_mission", "") - self._bank_retain_mission = self._config.get("bank_retain_mission") or None + self._bank_mission = cfg.get("bank_mission", "") + self._bank_retain_mission = cfg.get("bank_retain_mission") or None - # Tags self._retain_tags = _normalize_retain_tags( - self._config.get("retain_tags") - or os.environ.get("HINDSIGHT_RETAIN_TAGS", "") + cfg.get("retain_tags") or os.environ.get("HINDSIGHT_RETAIN_TAGS", "") ) self._tags = self._retain_tags or None self._observation_scopes = _normalize_observation_scopes( - self._config.get("observation_scopes") - or os.environ.get("HINDSIGHT_RETAIN_OBSERVATION_SCOPES", "") + cfg.get("observation_scopes") or os.environ.get("HINDSIGHT_RETAIN_OBSERVATION_SCOPES", "") ) - self._recall_tags = self._config.get("recall_tags") or None - self._recall_tags_match = self._config.get("recall_tags_match", "any") + self._recall_tags = cfg.get("recall_tags") or None + self._recall_tags_match = cfg.get("recall_tags_match", "any") self._retain_source = str( - self._config.get("retain_source") or os.environ.get("HINDSIGHT_RETAIN_SOURCE", _DEFAULT_RETAIN_SOURCE) + cfg.get("retain_source") or os.environ.get("HINDSIGHT_RETAIN_SOURCE", _DEFAULT_RETAIN_SOURCE) ).strip() self._retain_user_prefix = str( - self._config.get("retain_user_prefix") or os.environ.get("HINDSIGHT_RETAIN_USER_PREFIX", "User") + cfg.get("retain_user_prefix") or os.environ.get("HINDSIGHT_RETAIN_USER_PREFIX", "User") ).strip() or "User" self._retain_assistant_prefix = str( - self._config.get("retain_assistant_prefix") or os.environ.get("HINDSIGHT_RETAIN_ASSISTANT_PREFIX", "Assistant") + cfg.get("retain_assistant_prefix") or os.environ.get("HINDSIGHT_RETAIN_ASSISTANT_PREFIX", "Assistant") ).strip() or "Assistant" - # Retain controls - self._auto_retain = self._config.get("auto_retain", True) - self._retain_every_n_turns = max(1, int(self._config.get("retain_every_n_turns", 1))) - self._retain_context = self._config.get("retain_context", "conversation between Hermes Agent and the User") + self._auto_retain = cfg.get("auto_retain", True) + self._retain_every_n_turns = max(1, int(cfg.get("retain_every_n_turns", 1))) + self._retain_context = cfg.get("retain_context", _RETAIN_CONTEXT_DEFAULT) - # Recall controls - self._auto_recall = self._config.get("auto_recall", True) - self._recall_sync = bool(self._config.get("recall_sync", False)) - self._recall_max_tokens = int(self._config.get("recall_max_tokens", 4096)) - # Default narrows recall to observation-only; pass an explicit - # `recall_types` list in config.json to broaden (e.g. include - # "world" / "experience") or to disable the filter entirely. - configured_types = self._config.get("recall_types") + self._auto_recall = cfg.get("auto_recall", True) + self._recall_sync = bool(cfg.get("recall_sync", False)) + self._recall_max_tokens = int(cfg.get("recall_max_tokens", 4096)) + # None -> observation-only; a comma-separated string is accepted for + # parity with recall_tags; an explicit list broadens or disables the filter. + configured_types = cfg.get("recall_types") if configured_types is None: self._recall_types = ["observation"] elif isinstance(configured_types, str): - # Allow comma-separated strings for parity with recall_tags. self._recall_types = [t.strip() for t in configured_types.split(",") if t.strip()] else: self._recall_types = list(configured_types) or ["observation"] - self._recall_prompt_preamble = self._config.get("recall_prompt_preamble", "") - # On-by-default deterministic indicator: when auto-recall injects memory, - # Hermes emits a "👁️ Hindsight — recalled N memories" status line so the - # user SEES memory working, independent of whether the model mentions it. - # Off switch for customer-facing agents that shouldn't surface internals. - self._recall_indicator = bool(self._config.get("recall_indicator", True)) - # Companion retain indicator: "👁️ Hindsight — saving to memory…" emitted - # when a turn is dispatched to the writer. Same off switch rationale. - self._retain_indicator = bool(self._config.get("retain_indicator", True)) - self._recall_max_input_chars = int(self._config.get("recall_max_input_chars", 800)) - self._retain_async = self._config.get("retain_async", True) - self._prefetch_waits_for_retain = self._config.get("prefetch_waits_for_retain", True) - self._prefetch_retain_drain_timeout = float( - self._config.get("prefetch_retain_drain_timeout", 10.0) - ) + self._recall_prompt_preamble = cfg.get("recall_prompt_preamble", "") + # On by default so the user SEES memory working regardless of whether the + # model mentions it; off switch for customer-facing agents. + self._recall_indicator = bool(cfg.get("recall_indicator", True)) + self._retain_indicator = bool(cfg.get("retain_indicator", True)) + self._recall_max_input_chars = int(cfg.get("recall_max_input_chars", 800)) + self._retain_async = cfg.get("retain_async", True) + self._prefetch_waits_for_retain = cfg.get("prefetch_waits_for_retain", True) + self._prefetch_retain_drain_timeout = float(cfg.get("prefetch_retain_drain_timeout", 10.0)) - _client_version = "unknown" + client_version = "unknown" try: from importlib.metadata import version as pkg_version - _client_version = pkg_version("hindsight-client") + client_version = pkg_version("hindsight-client") except Exception: pass logger.info("Hindsight initialized: mode=%s, api_url=%s, bank=%s, budget=%s, memory_mode=%s, prefetch_method=%s, client=%s", - self._mode, self._api_url, self._bank_id, self._budget, self._memory_mode, self._prefetch_method, _client_version) + self._mode, self._api_url, self._bank_id, self._budget, self._memory_mode, self._prefetch_method, client_version) if self._bank_id_template: logger.debug("Hindsight bank resolved from template %r: profile=%s workspace=%s platform=%s user=%s -> bank=%s", self._bank_id_template, self._agent_identity, self._agent_workspace, @@ -1778,155 +983,128 @@ class HindsightMemoryProvider(MemoryProvider): self._retain_async, self._retain_context, self._recall_max_tokens, self._recall_max_input_chars, self._tags, self._recall_tags) - # For local mode, start the embedded daemon in the background so it - # doesn't block the chat. Redirect stdout/stderr to a log file to - # prevent rich startup output from spamming the terminal. if self._mode == "local_embedded": - # PostgreSQL's initdb refuses to run as root by design, so the - # embedded daemon can never initialize its data directory under - # root. Without this guard the daemon-start thread would fail, - # retry, and loop forever — each cycle reloading embedding models - # (~958MB RAM, ~33% CPU) with no user-visible error. Detect root - # up front and skip daemon startup with a clear message instead. - if hasattr(os, "geteuid") and os.geteuid() == 0: - msg = ( - "Hindsight local_embedded mode cannot run as root " - "(PostgreSQL initdb refuses root). Skipping the embedded " - "memory daemon. Run Hermes as a non-root user, or switch " - "to cloud / local_external mode via 'hermes memory setup'." - ) - logger.warning(msg) - # Surface to the terminal too — a daemon that never starts - # would otherwise fail silently and the user would only see - # Hermes get sluggish. (issue #13125) - try: - print(f" ⚠ {msg}", file=sys.stderr, flush=True) - except Exception: - pass - self._mode = "disabled" - return + self._start_embedded_daemon() - def _start_daemon(): - import traceback - log_dir = get_hermes_home() / "logs" - log_dir.mkdir(parents=True, exist_ok=True) - log_path = log_dir / "hindsight-embed.log" - try: - # Redirect the daemon manager's Rich console to our log file - # instead of stderr. This avoids global fd redirects that - # would capture output from other threads. - import hindsight_embed.daemon_embed_manager as dem - from rich.console import Console - dem.console = Console(file=open(log_path, "a", encoding="utf-8"), force_terminal=False) - - client = self._get_client() - profile = self._config.get("profile", "hermes") - - # Update the profile .env to match our current config so - # the daemon always starts with the right settings. - # If the config changed and the daemon is running, stop it. - profile_env = _embedded_profile_env_path(self._config) - expected_env = _build_embedded_profile_env(self._config) - saved = _load_simple_env(profile_env) - config_changed = saved != expected_env - - if config_changed: - profile_env = _materialize_embedded_profile_env(self._config) - if client._manager.is_running(profile): - with open(log_path, "a", encoding="utf-8") as f: - f.write("\n=== Config changed, restarting daemon ===\n") - client._manager.stop(profile) - - client._ensure_started() - with open(log_path, "a", encoding="utf-8") as f: - f.write("\n=== Daemon started successfully ===\n") - except Exception as e: - with open(log_path, "a", encoding="utf-8") as f: - f.write(f"\n=== Daemon startup failed: {e} ===\n") - traceback.print_exc(file=f) - - t = threading.Thread( - target=contextvars.copy_context().run, - args=(_start_daemon,), - daemon=True, - name="hindsight-daemon-start", + def _start_embedded_daemon(self) -> None: + """Start the embedded daemon on a background thread so it doesn't block + the chat; its Rich output goes to a log file, not the terminal.""" + # PostgreSQL's initdb refuses root, so the daemon can never initialize + # its data directory under root. Without this guard the start thread + # fails, retries and loops forever, reloading embedding models (~958MB + # RAM, ~33% CPU) with no user-visible error. + if hasattr(os, "geteuid") and os.geteuid() == 0: + msg = ( + "Hindsight local_embedded mode cannot run as root " + "(PostgreSQL initdb refuses root). Skipping the embedded " + "memory daemon. Run Hermes as a non-root user, or switch " + "to cloud / local_external mode via 'hermes memory setup'." ) - t.start() + logger.warning(msg) + # Also print: a daemon that never starts would otherwise fail + # silently and the user would only see Hermes get sluggish. + try: + print(f" ⚠ {msg}", file=sys.stderr, flush=True) + except Exception: + pass + self._mode = "disabled" + return + _context_thread(self._daemon_start_worker, "hindsight-daemon-start").start() + + def _daemon_start_worker(self) -> None: + import traceback + log_dir = get_hermes_home() / "logs" + log_dir.mkdir(parents=True, exist_ok=True) + log_path = log_dir / "hindsight-embed.log" + + def _log(text: str, exc: bool = False) -> None: + with open(log_path, "a", encoding="utf-8") as f: + f.write(text) + if exc: + traceback.print_exc(file=f) + + try: + # Point the daemon manager's Rich console at our log file rather + # than redirecting global fds (which would capture other threads). + import hindsight_embed.daemon_embed_manager as dem + from rich.console import Console + dem.console = Console(file=open(log_path, "a", encoding="utf-8"), force_terminal=False) + + client = self._get_client() + profile = self._config.get("profile", "hermes") + # Keep the profile .env in sync with our config; if it changed and + # the daemon is running, restart it with the new settings. + if _load_simple_env(_embedded_profile_env_path(self._config)) != _build_embedded_profile_env(self._config): + _materialize_embedded_profile_env(self._config) + if client._manager.is_running(profile): + _log("\n=== Config changed, restarting daemon ===\n") + client._manager.stop(profile) + client._ensure_started() + _log("\n=== Daemon started successfully ===\n") + except Exception as e: + _log(f"\n=== Daemon startup failed: {e} ===\n", exc=True) def system_prompt_block(self) -> str: - if self._memory_mode == "context": - return ( - f"# Hindsight Memory\n" - f"Active (context mode). Bank: {self._bank_id}, budget: {self._budget}.\n" - f"Relevant memories are automatically injected into context." - ) - if self._memory_mode == "tools": - return ( - f"# Hindsight Memory\n" - f"Active (tools mode). Bank: {self._bank_id}, budget: {self._budget}.\n" - f"Use hindsight_recall to search, hindsight_reflect for synthesis, " - f"hindsight_retain to store facts." - ) + mode = self._memory_mode if self._memory_mode in ("context", "tools") else "hybrid" + label = "" if mode == "hybrid" else f" ({mode} mode)" return ( f"# Hindsight Memory\n" - f"Active. Bank: {self._bank_id}, budget: {self._budget}.\n" - f"Relevant memories are automatically injected into context. " - f"Use hindsight_recall to search, hindsight_reflect for synthesis, " - f"hindsight_retain to store facts." + f"Active{label}. Bank: {self._bank_id}, budget: {self._budget}.\n" + f"{_SYSTEM_PROMPT_TAILS[mode]}" ) + # -- recall ------------------------------------------------------------------ + def _recall_disabled(self) -> bool: """Guards shared by the async and synchronous recall paths.""" - if self._memory_mode == "tools": - logger.debug("Prefetch: skipped (tools-only mode)") - return True - if not self._auto_recall: - logger.debug("Prefetch: skipped (auto_recall disabled)") - return True - if self._shutting_down.is_set(): - logger.debug("Prefetch: skipped (shutting down)") - return True + for skip, why in ( + (self._memory_mode == "tools", "tools-only mode"), + (not self._auto_recall, "auto_recall disabled"), + (self._shutting_down.is_set(), "shutting down"), + ): + if skip: + logger.debug("Prefetch: skipped (%s)", why) + return True return False - def _do_recall(self, query: str) -> _RecallResult: - """Run one recall/reflect for *query*. + def _recall_kwargs(self, query: str) -> dict: + kwargs: dict = { + "bank_id": self._bank_id, "query": query, + "budget": self._budget, "max_tokens": self._recall_max_tokens, + } + if self._recall_tags: + kwargs["tags"] = self._recall_tags + kwargs["tags_match"] = self._recall_tags_match + if self._recall_types: + kwargs["types"] = self._recall_types + return kwargs - Returns the formatted memory text plus the number of discrete memories - recalled (0 for a reflect synthesis or on error), so the deterministic - recall indicator can report an accurate count without re-parsing the - text. Shared by the background prefetch worker (``queue_prefetch``) and - the opt-in synchronous path (``prefetch`` when ``recall_sync`` is on). - """ - # Truncate query to max chars + def _do_recall(self, query: str) -> _RecallResult: + """One recall/reflect for *query*. Shared by the background prefetch + worker and the opt-in synchronous path (``recall_sync``).""" if self._recall_max_input_chars and len(query) > self._recall_max_input_chars: query = query[:self._recall_max_input_chars] try: if self._prefetch_method == "reflect": logger.debug("Recall: calling reflect (bank=%s, query_len=%d)", self._bank_id, len(query)) resp = self._run_hindsight_operation(lambda client: client.areflect(bank_id=self._bank_id, query=query, budget=self._budget)) - # Reflect synthesizes across many memories -> no discrete count. - return _RecallResult(resp.text or "", 0) - recall_kwargs: dict = { - "bank_id": self._bank_id, "query": query, - "budget": self._budget, "max_tokens": self._recall_max_tokens, - } - if self._recall_tags: - recall_kwargs["tags"] = self._recall_tags - recall_kwargs["tags_match"] = self._recall_tags_match - if self._recall_types: - recall_kwargs["types"] = self._recall_types + return _RecallResult(resp.text or "", 0) # synthesis -> no discrete count + recall_kwargs = self._recall_kwargs(query) logger.debug("Recall: calling recall (bank=%s, query_len=%d, budget=%s)", self._bank_id, len(query), self._budget) resp = self._run_hindsight_operation(lambda client: client.arecall(**recall_kwargs)) - num_results = len(resp.results) if resp.results else 0 - logger.debug("Recall: returned %d results", num_results) - text = "\n".join(f"- {r.text}" for r in resp.results if r.text) if resp.results else "" - return _RecallResult(text, num_results) + results = resp.results or [] + logger.debug("Recall: returned %d results", len(results)) + return _RecallResult("\n".join(f"- {r.text}" for r in results if r.text), len(results)) except Exception as e: logger.debug("Hindsight recall failed: %s", e, exc_info=True) return _RecallResult("", 0) - def _format_recall(self, result: str) -> str: + def _finish_prefetch(self, result: str, count: int) -> str: + """Record indicator state (cleared on empty turns so it never reports a + stale count) and format the injected block.""" + self._last_recall_returned = bool(result) + self._last_recall_count = count if result else 0 if not result: logger.debug("Prefetch: no results available") return "" @@ -1938,67 +1116,38 @@ class HindsightMemoryProvider(MemoryProvider): ) return f"{header}\n\n{result}" - def _record_recall_indicator(self, *, returned: bool, count: int) -> None: - """Track what the last prefetch injected, for recall_status(). - - Cleared to "nothing" on empty turns so the indicator never reports a - stale prior count. - """ - self._last_recall_returned = returned - self._last_recall_count = count if returned else 0 - def prefetch(self, query: str, *, session_id: str = "") -> str: # Opt-in: recall synchronously against the *current* message so the - # injected memories match this turn's query rather than the previous - # turn's queued recall. See NousResearch/hermes-agent#5820. + # injected memories match this turn's query, not the previous turn's. if self._recall_sync: if self._recall_disabled(): - self._record_recall_indicator(returned=False, count=0) - return "" + return self._finish_prefetch("", 0) recalled = self._do_recall(query) - self._record_recall_indicator(returned=bool(recalled.text), count=recalled.count) - return self._format_recall(recalled.text) - - # Default: return the result the background worker prefetched for the - # previous turn (cheap buffer read, capped join). + return self._finish_prefetch(recalled.text, recalled.count) + # Default: the background worker's result for the previous turn (capped join). if self._prefetch_thread and self._prefetch_thread.is_alive(): logger.debug("Prefetch: waiting for background thread to complete") self._prefetch_thread.join(timeout=3.0) with self._prefetch_lock: - result = self._prefetch_result - count = self._prefetch_count - self._prefetch_result = "" - self._prefetch_count = 0 - self._record_recall_indicator(returned=bool(result), count=count) - return self._format_recall(result) + result, count = self._prefetch_result, self._prefetch_count + self._prefetch_result, self._prefetch_count = "", 0 + return self._finish_prefetch(result, count) def recall_status(self) -> Optional[RecallStatus]: - """Report the count injected by the last prefetch (for the UI indicator). - - Returns ``None`` when nothing was injected this turn or the indicator - is turned off (``recall_indicator=false``), so customer-facing agents - can suppress the "recalled N memories" status line. - """ + """Count injected by the last prefetch, or None when nothing was injected + or ``recall_indicator=false`` (customer-facing agents).""" if not self._recall_indicator or not self._last_recall_returned: return None return RecallStatus(provider_label="Hindsight", count=self._last_recall_count, glyph=_HINDSIGHT_GLYPH) def queue_prefetch(self, query: str, *, session_id: str = "") -> None: - # In synchronous mode prefetch() does a live recall each turn, so - # there's nothing to prime in the background. - if self._recall_sync: - return - if self._recall_disabled(): + # Sync mode recalls live each turn — nothing to prime in the background. + if self._recall_sync or self._recall_disabled(): return def _run(): - # Ensure the just-completed turn's retain is recall-visible on the - # server before we recall, so the warmed context for the next turn - # includes it. This waits for the local writer queue to drain AND - # for the server-side async retain op(s) to complete (an explicit - # read-after-write signal), because async retain returns on - # acceptance rather than durability. Runs on the background prefetch - # thread, never the reply path, so it adds no response latency. + # Wait (bounded, off the reply path) for the just-completed turn's + # retain to be recall-visible so the warmed context includes it. if self._prefetch_waits_for_retain: self._wait_for_retains_drained(self._prefetch_retain_drain_timeout) recalled = self._do_recall(query) @@ -2007,29 +1156,17 @@ class HindsightMemoryProvider(MemoryProvider): self._prefetch_result = recalled.text self._prefetch_count = recalled.count - self._prefetch_thread = threading.Thread( - target=contextvars.copy_context().run, - args=(_run,), - daemon=True, - name="hindsight-prefetch", - ) + self._prefetch_thread = _context_thread(_run, "hindsight-prefetch") self._prefetch_thread.start() + # -- retain ------------------------------------------------------------------ + def _build_turn_messages(self, user_content: str, assistant_content: str) -> List[Dict[str, str]]: - # Hindsight receives this pair as one conversation turn, so both - # messages intentionally share the same turn-level event timestamp. + # One conversation turn -> both messages share the turn-level event timestamp. now = _event_timestamp() return [ - { - "role": "user", - "content": f"{self._retain_user_prefix}: {user_content}", - "timestamp": now, - }, - { - "role": "assistant", - "content": f"{self._retain_assistant_prefix}: {assistant_content}", - "timestamp": now, - }, + {"role": "user", "content": f"{self._retain_user_prefix}: {user_content}", "timestamp": now}, + {"role": "assistant", "content": f"{self._retain_assistant_prefix}: {assistant_content}", "timestamp": now}, ] def _build_metadata(self, *, message_count: int, turn_index: int) -> Dict[str, str]: @@ -2040,24 +1177,10 @@ class HindsightMemoryProvider(MemoryProvider): } if self._retain_source: metadata["source"] = self._retain_source - if self._session_id: - metadata["session_id"] = self._session_id - if self._platform: - metadata["platform"] = self._platform - if self._user_id: - metadata["user_id"] = self._user_id - if self._user_name: - metadata["user_name"] = self._user_name - if self._chat_id: - metadata["chat_id"] = self._chat_id - if self._chat_name: - metadata["chat_name"] = self._chat_name - if self._chat_type: - metadata["chat_type"] = self._chat_type - if self._thread_id: - metadata["thread_id"] = self._thread_id - if self._agent_identity: - metadata["agent_identity"] = self._agent_identity + for name in _METADATA_ATTRS: + value = getattr(self, f"_{name}") + if value: + metadata[name] = value return metadata def _build_retain_kwargs( @@ -2065,72 +1188,98 @@ class HindsightMemoryProvider(MemoryProvider): content: str, *, context: str | None = None, - document_id: str | None = None, metadata: Dict[str, str] | None = None, tags: List[str] | None = None, - retain_async: bool | None = None, occurred_at: str | None = None, + update_mode: str | None = None, ) -> Dict[str, Any]: - # The item-level timestamp is what the Hindsight server uses to resolve - # occurred_start/occurred_end (including relative phrases in content). - # An explicit occurred_at (from the retain tool) wins; otherwise default - # to the configured event clock so relative times still resolve (#93568). - kwargs: Dict[str, Any] = { - "bank_id": self._bank_id, + """Build one aretain_batch item. The item timestamp is what the server + uses to resolve occurred_start/occurred_end (incl. relative phrases in + content): an explicit occurred_at wins, else the configured event clock.""" + item: Dict[str, Any] = { "content": content, "metadata": metadata or self._build_metadata(message_count=1, turn_index=self._turn_index), - "timestamp": occurred_at.strip() if occurred_at and occurred_at.strip() else _event_timestamp(), + "timestamp": (occurred_at or "").strip() or _event_timestamp(), } if context is not None: - kwargs["context"] = context - if document_id: + item["context"] = context + if update_mode is not None: + item["update_mode"] = update_mode + merged_tags = _normalize_retain_tags(list(self._retain_tags) + _normalize_retain_tags(tags)) + if merged_tags: + item["tags"] = merged_tags + if self._observation_scopes: + item["observation_scopes"] = self._observation_scopes + return item + + def _lineage_tags(self) -> list[str]: + tags = [] + if self._session_id: + tags.append(f"session:{self._session_id}") + if self._parent_session_id: + tags.append(f"parent:{self._parent_session_id}") + return tags + + def _retain_batch(self, item: dict, *, bank_id: str, document_id: str | None = None, + retain_async: bool | None = None): + """Dispatch one item via aretain_batch (bank_id/document_id/retain_async are + call-level args, never item keys).""" + kwargs: Dict[str, Any] = {"bank_id": bank_id, "items": [item]} + if document_id is not None: kwargs["document_id"] = document_id if retain_async is not None: kwargs["retain_async"] = retain_async - merged_tags = _normalize_retain_tags(self._retain_tags) - for tag in _normalize_retain_tags(tags): - if tag not in merged_tags: - merged_tags.append(tag) - if merged_tags: - kwargs["tags"] = merged_tags - if self._observation_scopes: - kwargs["observation_scopes"] = self._observation_scopes - return kwargs + return self._run_hindsight_operation(lambda client: client.aretain_batch(**kwargs)) + + def _make_turn_retain_job(self, turns: list[str], *, document_id: str, update_mode: str | None, + label: str, track_ops: bool = True) -> Callable[[], None]: + """Build the writer job that ships *turns* as one document. Every input is + snapshotted NOW because the writer runs after later sync_turn() calls have + mutated _session_turns / _turn_index / _session_id.""" + content = "[" + ",".join(turns) + "]" + metadata = self._build_metadata(message_count=len(turns) * 2, turn_index=self._turn_index) + tags = self._lineage_tags() or None + bank_id, retain_async, retain_context = self._bank_id, self._retain_async, self._retain_context + + def _job() -> None: + item = self._build_retain_kwargs(content, context=retain_context, metadata=metadata, + tags=tags, update_mode=update_mode) + logger.debug("Hindsight %s: bank=%s, doc=%s, mode=%s, async=%s, content_len=%d, num_turns=%d", + label, bank_id, document_id, update_mode, retain_async, len(content), len(turns)) + resp = self._retain_batch(item, bank_id=bank_id, document_id=document_id, retain_async=retain_async) + # Async retains are only *accepted* here; track the op id(s) so the + # next-turn prefetch can wait for true server-side completion. + if retain_async and track_ops: + self._track_retain_ops(resp, bank_id) + logger.debug("Hindsight %s succeeded", label) + + return _job def sync_turn(self, user_content: str, assistant_content: str, *, session_id: str = "") -> None: - """Enqueue a retain for the current turn. Non-blocking. - - The actual aretain_batch runs on a single long-lived writer thread - that drains an in-memory queue. Once shutdown() has been called, - further sync_turn() calls are dropped — this prevents post-exit - retains from reaching aiohttp after interpreter shutdown begins. - """ + """Enqueue a retain for the current turn (non-blocking; runs on the writer + thread). Dropped once shutdown() has fired so post-exit retains never + reach aiohttp during interpreter teardown.""" if not self._auto_retain: logger.debug("sync_turn: skipped (auto_retain disabled)") return if self._shutting_down.is_set(): logger.debug("sync_turn: skipped (shutting down)") return - if session_id: self._session_id = str(session_id).strip() - turn = json.dumps(self._build_turn_messages(user_content, assistant_content), ensure_ascii=False) - self._session_turns.append(turn) + self._session_turns.append(json.dumps(self._build_turn_messages(user_content, assistant_content), ensure_ascii=False)) self._turn_counter += 1 self._turn_index = self._turn_counter - - if self._turn_counter % self._retain_every_n_turns != 0: + remainder = self._turn_counter % self._retain_every_n_turns + if remainder: logger.debug("sync_turn: buffered turn %d (will retain at turn %d)", - self._turn_counter, self._turn_counter + (self._retain_every_n_turns - self._turn_counter % self._retain_every_n_turns)) + self._turn_counter, self._turn_counter + (self._retain_every_n_turns - remainder)) return document_id, update_mode = self._resolve_retain_target(self._document_id) - - # On append-capable APIs each retain only needs to ship the turns - # accumulated since the last retain — the server appends them to the - # existing document. On legacy/overwrite APIs we must resend the whole - # session because each retain replaces the document. + # Append-capable APIs get only the delta since the last retain; legacy / + # overwrite APIs need the whole session because each retain replaces the document. if update_mode == "append": turns_to_retain = self._session_turns[self._last_retained_turn_count:] if not turns_to_retain: @@ -2138,77 +1287,25 @@ class HindsightMemoryProvider(MemoryProvider): return else: turns_to_retain = list(self._session_turns) - logger.debug("sync_turn: retaining %d/%d turns, payload %d chars", - len(turns_to_retain), len(self._session_turns), - sum(len(t) for t in turns_to_retain)) - content = "[" + ",".join(turns_to_retain) + "]" - - lineage_tags: list[str] = [] - if self._session_id: - lineage_tags.append(f"session:{self._session_id}") - if self._parent_session_id: - lineage_tags.append(f"parent:{self._parent_session_id}") - - # Snapshot the state needed for the retain. The writer may run after - # _session_turns / _turn_index are mutated by a later sync_turn(). - metadata_snapshot = self._build_metadata( - message_count=len(turns_to_retain) * 2, - turn_index=self._turn_index, - ) - num_turns = len(turns_to_retain) - bank_id = self._bank_id - retain_async_flag = self._retain_async - retain_context = self._retain_context - - def _do_retain() -> None: - item = self._build_retain_kwargs( - content, - context=retain_context, - metadata=metadata_snapshot, - tags=lineage_tags or None, - ) - item.pop("bank_id", None) - item.pop("retain_async", None) - if update_mode is not None: - item["update_mode"] = update_mode - logger.debug("Hindsight retain: bank=%s, doc=%s, mode=%s, async=%s, content_len=%d, num_turns=%d", - bank_id, document_id, update_mode, retain_async_flag, len(content), num_turns) - resp = self._run_hindsight_operation( - lambda client: client.aretain_batch( - bank_id=bank_id, - items=[item], - document_id=document_id, - retain_async=retain_async_flag, - ) - ) - # For async retains the write is only *accepted* here; track the - # returned operation id(s) so the next-turn prefetch can wait for - # true server-side completion (read-after-write) before recalling. - if retain_async_flag: - self._track_retain_ops(resp, bank_id) - logger.debug("Hindsight retain succeeded") + len(turns_to_retain), len(self._session_turns), sum(len(t) for t in turns_to_retain)) + job = self._make_turn_retain_job(turns_to_retain, document_id=document_id, + update_mode=update_mode, label="retain") self._ensure_writer() self._register_atexit() - # Deterministic "saving to memory" indicator — emitted the moment a - # real retain is dispatched (past every skip/buffer gate above), so it - # only fires on turns that actually persist. + # Emitted only once every skip/buffer gate above has passed, so it fires + # solely on turns that actually persist. self._emit_saving_indicator() - self._retain_queue.put(_do_retain) - # Advance the append watermark only after the delta is queued, so a - # later retain doesn't re-ship turns we've already handed to the writer. + self._retain_queue.put(job) + # Advance the watermark only after the delta is queued so a later retain + # doesn't re-ship turns already handed to the writer. if update_mode == "append": self._last_retained_turn_count = len(self._session_turns) def _emit_saving_indicator(self) -> None: - """Surface a model-independent "saving to memory" status line. - - Runs on the background sync worker (sync_turn's caller). No-ops when the - indicator is turned off (``retain_indicator=false``) or no status - channel was injected. Never raises — a status-line failure must not - derail the retain. - """ + """Model-independent "saving to memory" status line. No-op when + ``retain_indicator=false`` or no status channel; never raises.""" if not self._retain_indicator or self._status_callback is None: return try: @@ -2216,84 +1313,67 @@ class HindsightMemoryProvider(MemoryProvider): except Exception: logger.debug("Retain indicator emit failed (non-fatal)", exc_info=True) + # -- tools ------------------------------------------------------------------- + def get_tool_schemas(self) -> List[Dict[str, Any]]: if self._memory_mode == "context": return [] return [RETAIN_SCHEMA, RECALL_SCHEMA, REFLECT_SCHEMA] + def _tool_retain(self, args: dict) -> str: + content = args["content"] + context = args.get("context") + item = self._build_retain_kwargs(content, context=context, tags=args.get("tags"), + occurred_at=args.get("occurred_at")) + logger.debug("Tool hindsight_retain: bank=%s, content_len=%d, context=%s", + self._bank_id, len(content), context) + self._retain_batch(item, bank_id=self._bank_id) + logger.debug("Tool hindsight_retain: success") + return "Memory stored successfully." + + def _tool_recall(self, args: dict) -> str: + query = args["query"] + recall_kwargs = self._recall_kwargs(query) + logger.debug("Tool hindsight_recall: bank=%s, query_len=%d, budget=%s", + self._bank_id, len(query), self._budget) + resp = self._run_hindsight_operation(lambda client: client.arecall(**recall_kwargs)) + results = resp.results or [] + logger.debug("Tool hindsight_recall: %d results", len(results)) + if not results: + return "No relevant memories found." + return "\n".join(f"{i}. {r.text}" for i, r in enumerate(results, 1)) + + def _tool_reflect(self, args: dict) -> str: + query = args["query"] + logger.debug("Tool hindsight_reflect: bank=%s, query_len=%d, budget=%s", + self._bank_id, len(query), self._budget) + resp = self._run_hindsight_operation( + lambda client: client.areflect(bank_id=self._bank_id, query=query, budget=self._budget) + ) + logger.debug("Tool hindsight_reflect: response_len=%d", len(resp.text or "")) + return resp.text or "No relevant memories found." + + # tool name -> (required arg, handler) + _TOOL_HANDLERS = { + "hindsight_retain": ("content", _tool_retain), + "hindsight_recall": ("query", _tool_recall), + "hindsight_reflect": ("query", _tool_reflect), + } + def handle_tool_call(self, tool_name: str, args: dict, **kwargs) -> str: - if tool_name == "hindsight_retain": - content = args.get("content", "") - if not content: - return tool_error("Missing required parameter: content") - context = args.get("context") - try: - item = self._build_retain_kwargs( - content, - context=context, - tags=args.get("tags"), - occurred_at=args.get("occurred_at"), - ) - # aretain_batch takes bank_id/retain_async as call args, not item keys. - item.pop("bank_id", None) - item.pop("retain_async", None) - logger.debug("Tool hindsight_retain: bank=%s, content_len=%d, context=%s", - self._bank_id, len(content), context) - self._run_hindsight_operation( - lambda client: client.aretain_batch(bank_id=self._bank_id, items=[item]) - ) - logger.debug("Tool hindsight_retain: success") - return json.dumps({"result": "Memory stored successfully."}) - except Exception as e: - logger.warning("hindsight_retain failed: %s", e, exc_info=True) - return tool_error(f"Failed to store memory: {e}") + entry = self._TOOL_HANDLERS.get(tool_name) + if entry is None: + return tool_error(f"Unknown tool: {tool_name}") + required, handler = entry + if not args.get(required, ""): + return tool_error(f"Missing required parameter: {required}") + try: + return json.dumps({"result": handler(self, args)}) + except Exception as e: + logger.warning("%s failed: %s", tool_name, e, exc_info=True) + return tool_error(f"{_TOOL_ERRORS[tool_name]}: {e}") - elif tool_name == "hindsight_recall": - query = args.get("query", "") - if not query: - return tool_error("Missing required parameter: query") - try: - recall_kwargs: dict = { - "bank_id": self._bank_id, "query": query, "budget": self._budget, - "max_tokens": self._recall_max_tokens, - } - if self._recall_tags: - recall_kwargs["tags"] = self._recall_tags - recall_kwargs["tags_match"] = self._recall_tags_match - if self._recall_types: - recall_kwargs["types"] = self._recall_types - logger.debug("Tool hindsight_recall: bank=%s, query_len=%d, budget=%s", - self._bank_id, len(query), self._budget) - resp = self._run_hindsight_operation(lambda client: client.arecall(**recall_kwargs)) - num_results = len(resp.results) if resp.results else 0 - logger.debug("Tool hindsight_recall: %d results", num_results) - if not resp.results: - return json.dumps({"result": "No relevant memories found."}) - lines = [f"{i}. {r.text}" for i, r in enumerate(resp.results, 1)] - return json.dumps({"result": "\n".join(lines)}) - except Exception as e: - logger.warning("hindsight_recall failed: %s", e, exc_info=True) - return tool_error(f"Failed to search memory: {e}") - - elif tool_name == "hindsight_reflect": - query = args.get("query", "") - if not query: - return tool_error("Missing required parameter: query") - try: - logger.debug("Tool hindsight_reflect: bank=%s, query_len=%d, budget=%s", - self._bank_id, len(query), self._budget) - resp = self._run_hindsight_operation( - lambda client: client.areflect( - bank_id=self._bank_id, query=query, budget=self._budget - ) - ) - logger.debug("Tool hindsight_reflect: response_len=%d", len(resp.text or "")) - return json.dumps({"result": resp.text or "No relevant memories found."}) - except Exception as e: - logger.warning("hindsight_reflect failed: %s", e, exc_info=True) - return tool_error(f"Failed to reflect: {e}") - - return tool_error(f"Unknown tool: {tool_name}") + # -- session lifecycle ------------------------------------------------------- def on_session_switch( self, @@ -2303,118 +1383,58 @@ class HindsightMemoryProvider(MemoryProvider): reset: bool = False, **kwargs, ) -> None: - """Refresh cached per-session state when the agent rotates session_id. + """Refresh per-session state when the agent rotates session_id (/resume, + /branch, /reset, /new, context compression); otherwise writes land in + the previous session's document. - Fires on /resume, /branch, /reset, /new, and context compression. - Without this hook, initialize()-cached state (``_session_id``, - ``_document_id``, ``_session_turns``, ``_turn_counter``) would keep - pointing at the previous session and writes would land in the wrong - document. See hermes-agent#6672. - - Always update ``_session_id`` so metadata and tags on subsequent - retains reflect the active session. Always mint a fresh - ``_document_id`` so the new session's retain doesn't overwrite the - old session's document on vectorize-io/hindsight#1303. Always clear - the accumulated batch buffers (``_session_turns``, ``_turn_counter``, - ``_turn_index``) — even for /resume and /branch, the new session's - batching must start from zero so an in-flight retain doesn't flush - under the wrong ``_document_id``. - - Before clearing, flush any buffered turns under the *old* - ``_document_id``. Users who set ``retain_every_n_turns > 1`` would - otherwise silently lose whatever's in ``_session_turns`` at the - moment of switch — the same data-loss class as the shutdown race, - just at a different lifecycle event. - - Also wait for any in-flight prefetch from the old session and drop - its cached result; otherwise the new session's first ``prefetch()`` - could read stale recall text from before the switch. - - ``parent_session_id`` is recorded for lineage tags on future retains. - ``reset`` is accepted but not needed for Hindsight's state model — - buffer clearing is correct for every session switch, not only /reset. + Always: update ``_session_id`` (metadata/tags), mint a fresh + ``_document_id`` (so the new session can't overwrite the old document), + and clear the batch buffers — even for /resume and /branch, batching must + restart from zero so an in-flight retain doesn't flush under the wrong + document. Before clearing, flush buffered turns under the OLD ids + (``retain_every_n_turns > 1`` would otherwise silently lose them), and + join any in-flight prefetch and drop its cached result so the new + session's first ``prefetch()`` can't read stale recall. ``reset`` is + accepted but unneeded: buffer clearing is correct for every switch. """ new_id = str(new_session_id or "").strip() if not new_id: return - # 1. Flush any buffered turns under the OLD identifiers. Snapshot - # everything before mutating self._* so metadata + tags + doc_id - # all reference the old session consistently. + # 1. Flush buffered turns under the OLD identifiers, resolved BEFORE the + # rotation so the flush lands in the old session's document either way + # (legacy: per-process unique; >=0.5.0: stable session-scoped + append). if self._session_turns: - old_turns = list(self._session_turns) - old_session_id = self._session_id - old_parent_session_id = self._parent_session_id - old_turn_index = self._turn_index - old_metadata = self._build_metadata( - message_count=len(old_turns) * 2, - turn_index=old_turn_index, - ) - old_lineage_tags: list[str] = [] - if old_session_id: - old_lineage_tags.append(f"session:{old_session_id}") - if old_parent_session_id: - old_lineage_tags.append(f"parent:{old_parent_session_id}") - old_content = "[" + ",".join(old_turns) + "]" - # Resolve doc_id + update_mode against the OLD session BEFORE - # we rotate _session_id, so the flush lands in the old - # session's document either way (legacy: per-process unique; - # ≥0.5.0: stable session-scoped + append). - old_document_id, old_update_mode = self._resolve_retain_target( - self._document_id - ) + old_document_id, old_update_mode = self._resolve_retain_target(self._document_id) + job = self._make_turn_retain_job(list(self._session_turns), document_id=old_document_id, + update_mode=old_update_mode, label="flush-on-switch", + track_ops=False) def _flush(): try: - item = self._build_retain_kwargs( - old_content, - context=self._retain_context, - metadata=old_metadata, - tags=old_lineage_tags or None, - ) - item.pop("bank_id", None) - item.pop("retain_async", None) - if old_update_mode is not None: - item["update_mode"] = old_update_mode - logger.debug( - "Hindsight flush-on-switch: bank=%s, doc=%s, mode=%s, num_turns=%d", - self._bank_id, old_document_id, old_update_mode, len(old_turns), - ) - self._run_hindsight_operation( - lambda client: client.aretain_batch( - bank_id=self._bank_id, - items=[item], - document_id=old_document_id, - retain_async=self._retain_async, - ) - ) + job() except Exception as e: logger.warning("Hindsight flush-on-switch failed: %s", e, exc_info=True) - # Route the flush through the same writer queue sync_turn - # uses. That serializes it behind any still-queued retains - # from the old session (FIFO by document_id), avoids racing - # two threads on aretain_batch against the same document, and - # keeps shutdown's drain semantics intact. Skip enqueue if - # shutdown has already fired — the writer is draining/gone. + # Same writer queue as sync_turn: FIFO behind still-queued old-session + # retains, no two threads racing aretain_batch on one document, and + # shutdown's drain semantics intact. Skip once shutdown has fired. if not self._shutting_down.is_set(): self._ensure_writer() self._register_atexit() self._retain_queue.put(_flush) - # 2. Drain any in-flight prefetch from the old session and drop - # its cached result so the new session doesn't see stale recall. + # 2. Drain the old session's in-flight prefetch and drop its result. if self._prefetch_thread and self._prefetch_thread.is_alive(): self._prefetch_thread.join(timeout=3.0) with self._prefetch_lock: self._prefetch_result = "" - # 3. Now rotate to the new session. + # 3. Rotate to the new session. if parent_session_id: self._parent_session_id = str(parent_session_id).strip() self._session_id = new_id - start_ts = datetime.now().strftime("%Y%m%d_%H%M%S_%f") - self._document_id = f"{self._session_id}-{start_ts}" + self._document_id = _mint_document_id(new_id) self._session_turns = [] self._turn_counter = 0 self._turn_index = 0 @@ -2426,12 +1446,10 @@ class HindsightMemoryProvider(MemoryProvider): def shutdown(self) -> None: logger.debug("Hindsight shutdown: stopping writer + waiting for background threads") - # Stop accepting new retain jobs first so anyone still calling - # sync_turn() during teardown is dropped, not enqueued. + # Stop accepting retain jobs first so late sync_turn() calls are dropped. self._shutting_down.set() - # Drain the writer: it will finish in-flight work, then exit on - # the sentinel. Bounded join keeps shutdown predictable even if - # the daemon is wedged. + # The writer finishes in-flight work then exits on the sentinel; the + # bounded join keeps shutdown predictable even if the daemon is wedged. writer = self._writer_thread if writer is not None and writer.is_alive(): try: @@ -2440,22 +1458,17 @@ class HindsightMemoryProvider(MemoryProvider): pass writer.join(timeout=10.0) if writer.is_alive(): - logger.warning( - "Hindsight writer did not stop within 10s; " - "abandoning %d pending retain(s)", - self._retain_queue.qsize(), - ) + logger.warning("Hindsight writer did not stop within 10s; abandoning %d pending retain(s)", + self._retain_queue.qsize()) if self._prefetch_thread and self._prefetch_thread.is_alive(): self._prefetch_thread.join(timeout=5.0) if self._client is not None: try: if self._mode == "local_embedded": - # HindsightEmbedded.close() delegates to its sync client.close(). - # When Hermes created/used that client on the shared async loop, - # closing it from this thread can raise "attached to a different - # loop" before aiohttp releases the session. Close the embedded - # inner async client on the shared loop first, then let the - # wrapper clean up daemon/UI bookkeeping. + # HindsightEmbedded.close() closes its sync client from this + # thread, which can raise "attached to a different loop" before + # aiohttp releases the session. Close the inner async client on + # the shared loop first, then let the wrapper clean up bookkeeping. inner_client = getattr(self._client, "_client", None) if inner_client is not None and hasattr(inner_client, "aclose"): _run_sync(inner_client.aclose()) @@ -2472,17 +1485,11 @@ class HindsightMemoryProvider(MemoryProvider): except Exception: pass self._client = None - # The module-global background event loop (_loop / _loop_thread) - # is intentionally NOT stopped here. It is shared across every - # HindsightMemoryProvider instance in the process — the plugin - # loader creates a new provider per AIAgent, and the gateway - # creates one AIAgent per concurrent chat session. Stopping the - # loop from one provider's shutdown() strands the aiohttp - # ClientSession + TCPConnector owned by every sibling provider - # on a dead loop, which surfaces as the "Unclosed client session" - # / "Unclosed connector" warnings reported in #11923. The loop - # runs on a daemon thread and is reclaimed on process exit; - # per-session cleanup happens via self._client.aclose() above. + # The module-global loop (_loop / _loop_thread) is intentionally NOT + # stopped: it's shared by every provider in the process (one per AIAgent, + # one AIAgent per gateway chat session). Stopping it would strand sibling + # providers' aiohttp sessions on a dead loop ("Unclosed client session"). + # It runs on a daemon thread and is reclaimed at process exit. def register(ctx) -> None: diff --git a/plugins/memory/hindsight/embedded.py b/plugins/memory/hindsight/embedded.py new file mode 100644 index 0000000000..3283dad9be --- /dev/null +++ b/plugins/memory/hindsight/embedded.py @@ -0,0 +1,195 @@ +"""Local-embedded Hindsight runtime helpers: import probe, install hint, the +per-profile env file the standalone ``hindsight-embed`` daemon consumes, and the +daemon health-grace env export.""" + +from __future__ import annotations + +import importlib +import logging +import os +import sys +from pathlib import Path +from typing import Any + +from agent.secret_scope import get_secret + +from .settings import _DEFAULT_IDLE_TIMEOUT, _daemon_llm_provider, _parse_int_setting + +logger = logging.getLogger(__name__.rpartition(".")[0]) + +# Read by hindsight_embed.daemon_embed_manager AT IMPORT TIME (module-level +# constant): how long to wait for a slow /health before declaring the daemon +# stale and killing it. On resource-contended hosts a busy daemon can exceed +# the upstream 2s check and get needlessly restarted, so it's plugin config. +_PORT_HEALTH_GRACE_ENV = "HINDSIGHT_EMBED_PORT_HEALTH_GRACE_TIMEOUT" + +# Markers of a stale embedded-daemon connection (the client is recreated and +# the operation retried once). +_RETRIABLE_CONNECTION_MARKERS = ( + "cannot connect to host", + "connection refused", + "connect call failed", + "clientconnectorerror", +) + + +def _export_port_health_grace_timeout(config: dict[str, Any]) -> None: + """Export the daemon health grace timeout to the process env. + + Must run BEFORE ``hindsight_embed.daemon_embed_manager`` is imported. Only + set when the user configured a value; ``setdefault`` so an explicit env + override always wins. + """ + raw = config.get("port_health_grace_timeout") + if raw is None or raw == "": + return + try: + seconds = float(raw) + except (TypeError, ValueError): + logger.warning("Invalid Hindsight port_health_grace_timeout %r; ignoring.", raw) + return + if seconds < 0: + logger.warning("Negative Hindsight port_health_grace_timeout %r; ignoring.", raw) + return + os.environ.setdefault(_PORT_HEALTH_GRACE_ENV, repr(seconds)) + + +def _check_local_runtime() -> tuple[bool, str | None]: + """Return whether the local embedded Hindsight stack imports cleanly. + + On older CPUs NumPy can raise at import before the daemon starts; report + "unavailable" so Hermes degrades instead of retrying a broken backend. + ``sentence_transformers`` is imported too: ``hindsight``/``hindsight_embed`` + import fine even when the embedding stack is broken, and without this the + probe (and ``hermes memory status``) would stay green while the daemon + aborts on every retain/recall. + """ + try: + importlib.import_module("hindsight") + importlib.import_module("hindsight_embed.daemon_embed_manager") + importlib.import_module("sentence_transformers") + return True, None + except Exception as exc: + return False, str(exc) + + +def _local_runtime_hint(reason: str | None) -> str: + """Install guidance when the local_embedded runtime is missing. + + The top-level ``hindsight`` module ships only with ``hindsight-all``; + ``plugin.yaml`` declares just ``hindsight-client`` (cloud/local_external), + so a hand-written config, the legacy ``"mode": "local"`` alias, or a + restored backup hits ``ModuleNotFoundError: No module named 'hindsight'``. + """ + text = (reason or "").lower() + if "no module named" in text and ("hindsight'" in text or 'hindsight"' in text + or "hindsight_embed" in text): + return ( + f" Install the embedded runtime with: uv pip install --python " + f"{sys.executable} hindsight-all — or run 'hermes memory setup'. " + "(local_embedded needs the 'hindsight-all' package, which provides the " + "top-level 'hindsight' module; 'hindsight-client' alone only covers " + "cloud / local_external.)" + ) + return "" + + +def _load_simple_env(path) -> dict[str, str]: + """Parse a KEY=VALUE env file, ignoring comments and blank lines. + + utf-8-sig, not utf-8: also used on the Hermes .env during post_setup, where + a Notepad BOM would otherwise stick to the first key. + """ + if not path.exists(): + return {} + values: dict[str, str] = {} + for line in path.read_text(encoding="utf-8-sig", errors="replace").splitlines(): + if not line or line.startswith("#") or "=" not in line: + continue + key, value = line.split("=", 1) + values[key.strip()] = value.strip() + return values + + +def _embedded_profile_env_path(config: dict[str, Any]) -> Path: + profile = str(config.get("profile", "hermes") or "hermes") + return Path.home() / ".hindsight" / "profiles" / f"{profile}.env" + + +def _build_embedded_profile_env(config: dict[str, Any], *, llm_api_key: str | None = None) -> dict[str, str]: + """Build the profile-scoped env that standalone hindsight-embed consumes.""" + if llm_api_key is None: + llm_api_key = ( + config.get("llmApiKey") + or config.get("llm_api_key") + or get_secret("HINDSIGHT_LLM_API_KEY", "") + ) + env_values = { + "HINDSIGHT_API_LLM_PROVIDER": str(_daemon_llm_provider(config.get("llm_provider", ""))), + "HINDSIGHT_API_LLM_API_KEY": str(llm_api_key or ""), + "HINDSIGHT_API_LLM_MODEL": str(config.get("llm_model", "")), + "HINDSIGHT_API_LOG_LEVEL": "info", + } + base_url = config.get("llm_base_url") or os.environ.get("HINDSIGHT_API_LLM_BASE_URL", "") + if base_url: + env_values["HINDSIGHT_API_LLM_BASE_URL"] = str(base_url) + idle_timeout = config.get("idle_timeout") + if idle_timeout is None: + idle_timeout = os.environ.get("HINDSIGHT_IDLE_TIMEOUT") + if idle_timeout is not None and idle_timeout != "": + env_values["HINDSIGHT_EMBED_DAEMON_IDLE_TIMEOUT"] = str( + _parse_int_setting(idle_timeout, _DEFAULT_IDLE_TIMEOUT) + ) + return env_values + + +def _secure_write_profile_env(profile_env: Path, content: str) -> None: + """Create/overwrite *profile_env* owner-only (0600). The file carries the + daemon's plaintext LLM API key, so a pre-existing file is tightened BEFORE + the new secret bytes are written.""" + if profile_env.exists(): + try: + os.chmod(profile_env, 0o600) + except OSError: + pass + fd = os.open(str(profile_env), os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600) + with os.fdopen(fd, "w", encoding="utf-8") as fh: + fh.write(content) + + +def _validate_profile_env_permissions(profile_env: Path) -> None: + """Post-write check: the secret file must be owner-only on POSIX (Windows + ACLs aren't modelled by mode bits, so skipped there).""" + if os.name != "posix": + return + import stat + + mode = stat.S_IMODE(profile_env.stat().st_mode) + if mode != 0o600: + try: + os.chmod(profile_env, 0o600) + except OSError: + pass + if stat.S_IMODE(profile_env.stat().st_mode) != 0o600: + raise PermissionError( + f"Embedded Hindsight profile environment is not owner-only: {profile_env}" + ) + + +def _materialize_embedded_profile_env(config: dict[str, Any], *, llm_api_key: str | None = None) -> Path: + """Write the profile env file; never leave a plaintext key behind in a file + whose permissions could not be verified.""" + profile_env = _embedded_profile_env_path(config) + profile_env.parent.mkdir(parents=True, exist_ok=True) + env_values = _build_embedded_profile_env(config, llm_api_key=llm_api_key) + content = "".join(f"{key}={value}\n" for key, value in env_values.items()) + try: + _secure_write_profile_env(profile_env, content) + _validate_profile_env_permissions(profile_env) + except BaseException: + try: + profile_env.unlink() + except OSError: + pass + raise + return profile_env diff --git a/plugins/memory/hindsight/settings.py b/plugins/memory/hindsight/settings.py new file mode 100644 index 0000000000..7c436d692d --- /dev/null +++ b/plugins/memory/hindsight/settings.py @@ -0,0 +1,158 @@ +"""Hindsight plugin constants and pure config normalizers (no I/O, no origin imports).""" + +from __future__ import annotations + +import json +import logging +from typing import Any, List + +# Log under the plugin package's own logger name (loader-path independent). +logger = logging.getLogger(__name__.rpartition(".")[0]) + +_DEFAULT_API_URL = "https://api.hindsight.vectorize.io" +_DEFAULT_LOCAL_URL = "http://localhost:8888" +# Keep in sync with tools/lazy_deps.py ("memory.hindsight") and plugin.yaml. +_MIN_CLIENT_VERSION = "0.6.1" +_DEFAULT_TIMEOUT = 120 # seconds — cloud API can take 30-40s per request +_DEFAULT_IDLE_TIMEOUT = 300 # seconds — Hindsight embedded daemon default +# ``metadata.source`` stamped on retained memories — OPT-IN, empty by default: +# AGENTS.md forbids on-by-default third-party attribution tags. Set via the +# ``retain_source`` config key or HINDSIGHT_RETAIN_SOURCE. +_DEFAULT_RETAIN_SOURCE = "" +# Hindsight brand mark (eye ringed by graph nodes) for the recall/retain indicators. +_HINDSIGHT_GLYPH = "👁️" +# Hindsight 0.5.0 added ``update_mode='append'`` on retain. Without it, reusing a +# stable session-scoped document_id silently overwrites prior turns server-side, +# so older APIs keep the per-process unique document_id fallback. +_MIN_VERSION_FOR_UPDATE_MODE_APPEND = "0.5.0" +_VALID_BUDGETS = {"low", "mid", "high"} +_PROVIDER_DEFAULT_MODELS = { + "openai": "gpt-4o-mini", + "anthropic": "claude-haiku-4-5", + "gemini": "gemini-3.6-flash", + "groq": "openai/gpt-oss-120b", + "openrouter": "qwen/qwen3.5-9b", + "minimax": "MiniMax-M2.7", + "ollama": "gemma3:12b", + "lmstudio": "local-model", + "openai_compatible": "your-model-name", +} +# The embedded daemon speaks OpenAI wire format for these providers. +_OPENAI_WIRE_PROVIDERS = {"openai_compatible", "openrouter"} +_OBSERVATION_SCOPE_KEYWORDS = {"per_tag", "combined", "all_combinations"} + + +def _parse_int_setting(value: Any, default: int) -> int: + """Parse an integer config/env value, falling back on invalid input.""" + if value is None or value == "": + return default + try: + return int(value) + except (TypeError, ValueError): + logger.warning("Invalid integer Hindsight setting %r; using default %s", value, default) + return default + + +def _daemon_llm_provider(provider: str) -> str: + return "openai" if provider in _OPENAI_WIRE_PROVIDERS else provider + + +def _normalize_retain_tags(value: Any) -> List[str]: + """Normalize tag config/tool values to a deduplicated list of strings.""" + if value is None: + return [] + if isinstance(value, list): + raw_items = value + elif isinstance(value, str): + text = value.strip() + if not text: + return [] + parsed = None + if text.startswith("["): + try: + parsed = json.loads(text) + except Exception: + parsed = None + raw_items = parsed if isinstance(parsed, list) else text.split(",") + else: + raw_items = [value] + normalized: list[str] = [] + for item in raw_items: + tag = str(item).strip() + if tag and tag not in normalized: + normalized.append(tag) + return normalized + + +def _normalize_observation_scopes(value: Any) -> Any: + """Normalize an observation_scopes value to a Hindsight-accepted form. + + Returns ``None`` (nothing configured; Hindsight applies its ``combined`` + default), a keyword string, or ``list[list[str]]`` (one inner list per + consolidation pass). Accepts a keyword, a JSON-encoded list, a flat list of + tags (one scope), or a list of tag-lists. Anything unrecognized yields + ``None`` so we never send an invalid payload. + """ + if isinstance(value, str): + text = value.strip() + if text in _OBSERVATION_SCOPE_KEYWORDS: + return text + if text.startswith("["): + try: + return _normalize_observation_scopes(json.loads(text)) + except Exception: + return None + return None + if isinstance(value, (list, tuple)): + if all(isinstance(entry, str) for entry in value): + inner = [entry.strip() for entry in value if entry.strip()] + return [inner] if inner else None + scopes: list[list[str]] = [] + for entry in value: + if isinstance(entry, (list, tuple)): + inner = [str(tag).strip() for tag in entry if str(tag).strip()] + if inner: + scopes.append(inner) + elif isinstance(entry, str) and entry.strip(): + scopes.append([entry.strip()]) + return scopes or None + return None + + +def _sanitize_bank_segment(value: str) -> str: + """Make a bank_id placeholder URL/filesystem safe: non ``[A-Za-z0-9_-]`` + runs become a single dash; leading/trailing dashes and underscores are stripped.""" + if not value: + return "" + out = [] + prev_dash = False + for ch in str(value): + if ch.isalnum() or ch in "-_": + out.append(ch) + prev_dash = False + elif not prev_dash: + out.append("-") + prev_dash = True + return "".join(out).strip("-_") + + +def _resolve_bank_id_template(template: str, fallback: str, **placeholders: str) -> str: + """Render a bank_id template ({profile}, {workspace}, {platform}, {user}, + {session}); each placeholder is sanitized first. Empty placeholders render + as "" and the dash/underscore runs they leave are collapsed, e.g. + ``hermes-{user}`` with no user becomes ``hermes``. Empty template or an + invalid placeholder falls back to *fallback*.""" + if not template: + return fallback + sanitized = {k: _sanitize_bank_segment(v) for k, v in placeholders.items()} + try: + rendered = template.format(**sanitized) + except (KeyError, IndexError) as exc: + logger.warning("Invalid bank_id_template %r: %s — using fallback %r", + template, exc, fallback) + return fallback + while "--" in rendered: + rendered = rendered.replace("--", "-") + while "__" in rendered: + rendered = rendered.replace("__", "_") + return rendered.strip("-_") or fallback diff --git a/plugins/memory/hindsight/setup.py b/plugins/memory/hindsight/setup.py new file mode 100644 index 0000000000..f5e0bddc6a --- /dev/null +++ b/plugins/memory/hindsight/setup.py @@ -0,0 +1,215 @@ +"""`hermes memory setup` wizard for the Hindsight provider (``post_setup``).""" + +from __future__ import annotations + +import json +import sys +from pathlib import Path + +from agent.secret_scope import get_secret +from hermes_cli.secret_prompt import masked_secret_prompt + +from . import templates as _hs_templates +from .embedded import ( + _embedded_profile_env_path, + _load_simple_env, + _materialize_embedded_profile_env, +) +from .settings import ( + _DEFAULT_API_URL, + _DEFAULT_IDLE_TIMEOUT, + _DEFAULT_LOCAL_URL, + _DEFAULT_TIMEOUT, + _MIN_CLIENT_VERSION, + _PROVIDER_DEFAULT_MODELS, +) + +_MODE_VALUES = ["cloud", "local_embedded", "local_external"] +_MODE_ITEMS = [ + ("Cloud", "Hindsight Cloud API (lightweight, just needs an API key)"), + ("Local Embedded", "Run Hindsight locally (downloads ~200MB, needs LLM key)"), + ("Local External", "Connect to an existing Hindsight instance"), +] + + +def _secret_prompt(label: str) -> str: + """Masked prompt on a TTY; plain readline when stdin is piped.""" + sys.stdout.write(label) + sys.stdout.flush() + return masked_secret_prompt("") if sys.stdin.isatty() else sys.stdin.readline().strip() + + +def _select(title: str, items: list, values: list, current) -> str | None: + """Curses pick from *values*, defaulting to *current*; None when cancelled.""" + from hermes_cli.memory_setup import _CANCELLED, _curses_select, _print_cancelled_setup + + default = values.index(current) if current in values else 0 + idx = _curses_select(title, items, default=default, cancel_returns=_CANCELLED) + if idx == _CANCELLED: + _print_cancelled_setup() + return None + return values[idx] + + +def _write_env(env_path: Path, env_writes: dict) -> None: + """Update keys in place (BOM-tolerant read so a Notepad BOM can't glue + U+FEFF onto the first key and cause a duplicate line), append the rest.""" + env_path.parent.mkdir(parents=True, exist_ok=True) + existing = env_path.read_text(encoding="utf-8-sig").splitlines() if env_path.exists() else [] + updated = set() + new_lines = [] + for line in existing: + key = line.split("=", 1)[0].strip() if "=" in line and not line.startswith("#") else None + if key in env_writes: + new_lines.append(f"{key}={env_writes[key]}") + updated.add(key) + else: + new_lines.append(line) + new_lines.extend(f"{k}={v}" for k, v in env_writes.items() if k not in updated) + env_path.write_text("\n".join(new_lines) + "\n", encoding="utf-8") + + +def _offer_starter_template(mode: str, provider_config: dict, env_writes: dict) -> None: + """Seed the bank with a Hermes starter template (best-effort).""" + import os + + from hermes_cli.memory_setup import _CANCELLED, _curses_select + + default_url = _DEFAULT_LOCAL_URL if mode == "local_external" else _DEFAULT_API_URL + _hs_templates.run_template_step( + api_url=provider_config.get("api_url") or default_url, + bank_id=provider_config.get("bank_id", "hermes"), + api_key=env_writes.get("HINDSIGHT_API_KEY") or os.environ.get("HINDSIGHT_API_KEY", "") or None, + select=_curses_select, + cancelled=_CANCELLED, + ) + + +def run_setup(provider, hermes_home: str, config: dict) -> None: + """Interactive wizard — installs only the deps the selected mode needs.""" + from hermes_cli.config import save_config + + from . import _load_config + + print("\n Configuring Hindsight memory:\n") + + existing_config = provider._config if isinstance(provider._config, dict) else _load_config() + if not isinstance(existing_config, dict): + existing_config = {} + + mode = _select(" Select mode", _MODE_ITEMS, _MODE_VALUES, existing_config.get("mode")) + if mode is None: + return + provider_config: dict = dict(existing_config, mode=mode) + env_writes: dict = {} + + llm_provider = "" + if mode == "local_embedded": + providers = list(_PROVIDER_DEFAULT_MODELS) + llm_items = [(p, f"default model: {_PROVIDER_DEFAULT_MODELS[p]}") for p in providers] + llm_provider = _select(" Select LLM provider", llm_items, providers, provider_config.get("llm_provider")) + if llm_provider is None: + return + provider_config["llm_provider"] = llm_provider + + print("\n Checking dependencies...") + # Environment-aware install: sealed hosted venvs redirect to the durable + # data-volume target instead of writing to /opt/hermes. + from tools.lazy_deps import install_specs + + deps = ["hindsight-all"] if mode == "local_embedded" else [f"hindsight-client>={_MIN_CLIENT_VERSION}"] + outcome = install_specs(deps, timeout=120) + if outcome.ok: + print(" ✓ Dependencies up to date") + elif outcome.blocked: + print(f" ⚠ Cannot install dependencies: {outcome.reason}") + else: + print(f" ⚠ Install failed:\n{(outcome.stderr or '').strip()}") + print(f" Run manually: uv pip install --python {sys.executable} {' '.join(deps)}") + + if mode == "cloud": + print("\n Get your API key at https://ui.hindsight.vectorize.io\n") + existing_key = get_secret("HINDSIGHT_API_KEY", "") or "" + if existing_key: + masked = f"...{existing_key[-4:]}" if len(existing_key) > 4 else "set" + api_key = _secret_prompt(f" API key (current: {masked}, blank to keep): ") + else: + api_key = _secret_prompt(" API key: ") + if api_key: + env_writes["HINDSIGHT_API_KEY"] = api_key + val = input(f" API URL [{_DEFAULT_API_URL}]: ").strip() + if val: + provider_config["api_url"] = val + + elif mode == "local_external": + val = input(f" Hindsight API URL [{_DEFAULT_LOCAL_URL}]: ").strip() + provider_config["api_url"] = val or _DEFAULT_LOCAL_URL + api_key = _secret_prompt(" API key (optional, blank to skip): ") + if api_key: + env_writes["HINDSIGHT_API_KEY"] = api_key + + else: # local_embedded + if llm_provider == "openai_compatible": + existing_base_url = provider_config.get("llm_base_url", "") + prompt = " LLM endpoint URL (e.g. http://192.168.1.10:8080/v1)" + if existing_base_url: + prompt += f" [{existing_base_url}]" + val = input(prompt + ": ").strip() + if val: + provider_config["llm_base_url"] = val + elif llm_provider == "openrouter": + provider_config["llm_base_url"] = "https://openrouter.ai/api/v1" + + current_model = provider_config.get("llm_model") or _PROVIDER_DEFAULT_MODELS.get(llm_provider, "gpt-4o-mini") + val = input(f" LLM model [{current_model}]: ").strip() + provider_config["llm_model"] = val or current_model + + llm_key = _secret_prompt(" LLM API key: ") + env_writes["HINDSIGHT_LLM_API_KEY"] = ( + llm_key or _load_simple_env(Path(hermes_home) / ".env").get("HINDSIGHT_LLM_API_KEY", "") + ) + + provider_config.setdefault("bank_id", "hermes") + provider_config.setdefault("recall_budget", "mid") + # Preserve explicit 0 timeouts instead of treating them as blank. + timeout_val = provider_config.get("timeout") + if timeout_val is None: + timeout_val = _DEFAULT_TIMEOUT + provider_config["timeout"] = timeout_val + env_writes["HINDSIGHT_TIMEOUT"] = str(timeout_val) + if mode == "local_embedded": + idle_timeout_val = provider_config.get("idle_timeout") + if idle_timeout_val is None: + idle_timeout_val = _DEFAULT_IDLE_TIMEOUT + provider_config["idle_timeout"] = idle_timeout_val + env_writes["HINDSIGHT_IDLE_TIMEOUT"] = str(idle_timeout_val) + config["memory"]["provider"] = "hindsight" + save_config(config) + provider.save_config(provider_config, hermes_home) + if env_writes: + _write_env(Path(hermes_home) / ".env", env_writes) + + # Starter template only for cloud / local_external — the API is reachable + # now; local_embedded's daemon isn't up during setup. + if _hs_templates.supported_for_mode(mode): + _offer_starter_template(mode, provider_config, env_writes) + + if mode == "local_embedded": + materialized_config = dict(provider_config) + try: + materialized_config = json.loads( + (Path(hermes_home) / "hindsight" / "config.json").read_text(encoding="utf-8") + ) + except Exception: + pass + llm_api_key = ( + env_writes.get("HINDSIGHT_LLM_API_KEY", "") + or _load_simple_env(Path(hermes_home) / ".env").get("HINDSIGHT_LLM_API_KEY", "") + or _load_simple_env(_embedded_profile_env_path(materialized_config)).get("HINDSIGHT_API_LLM_API_KEY", "") + ) + _materialize_embedded_profile_env(materialized_config, llm_api_key=llm_api_key or None) + + print(f"\n ✓ Hindsight memory configured ({mode} mode)") + if env_writes: + print(" API keys saved to .env") + print("\n Start a new session to activate.\n") diff --git a/tests/plugins/memory/test_hindsight_env_perms.py b/tests/plugins/memory/test_hindsight_env_perms.py index f444d4439b..f18143ba7c 100644 --- a/tests/plugins/memory/test_hindsight_env_perms.py +++ b/tests/plugins/memory/test_hindsight_env_perms.py @@ -65,12 +65,12 @@ def test_rewrite_tightens_existing_world_readable_profile_env(): def test_secret_file_removed_when_permission_validation_fails(monkeypatch): """If the post-write permission check cannot verify 0600, the plaintext key file must not be left behind.""" - import plugins.memory.hindsight as hs + import plugins.memory.hindsight.embedded as hs_embedded def _fail_validation(profile_env): raise PermissionError(f"not owner-only: {profile_env}") - monkeypatch.setattr(hs, "_validate_profile_env_permissions", _fail_validation) + monkeypatch.setattr(hs_embedded, "_validate_profile_env_permissions", _fail_validation) with pytest.raises(PermissionError): _materialize_embedded_profile_env(_CONFIG, llm_api_key="sk-doomed") diff --git a/tests/plugins/memory/test_hindsight_provider.py b/tests/plugins/memory/test_hindsight_provider.py index b48e2eb383..23c05e29bd 100644 --- a/tests/plugins/memory/test_hindsight_provider.py +++ b/tests/plugins/memory/test_hindsight_provider.py @@ -1413,7 +1413,7 @@ class TestAvailability: ) monkeypatch.setattr( - "plugins.memory.hindsight.importlib.import_module", + "importlib.import_module", _raise, ) p = HindsightMemoryProvider() @@ -1432,7 +1432,7 @@ class TestAvailability: raise RuntimeError("x86_64-v2 unsupported") monkeypatch.setattr( - "plugins.memory.hindsight.importlib.import_module", + "importlib.import_module", _raise, )