From 1d5d0594103fbd2d86e2533a373a945e784a208b Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Mon, 7 Sep 2026 03:22:27 -0700 Subject: [PATCH] fix(gateway): stop time-triggered conversation rotation --- gateway/__init__.py | 4 +- gateway/config.py | 41 +------ gateway/config_env.py | 9 +- gateway/config_loader.py | 1 - gateway/run.py | 9 +- gateway/run_agent_cache.py | 32 +----- gateway/run_startup.py | 2 +- gateway/run_turn.py | 31 +---- gateway/run_watchers.py | 89 +-------------- gateway/session.py | 6 +- gateway/session_lifecycle.py | 108 +----------------- gateway/session_persistence.py | 2 +- gateway/session_recovery.py | 13 +-- hermes_cli/config.py | 2 +- hermes_cli/setup.py | 51 --------- hermes_cli/setup_quick.py | 1 - .../gateway/test_session_time_persistence.py | 41 +++++++ 17 files changed, 75 insertions(+), 367 deletions(-) create mode 100644 tests/gateway/test_session_time_persistence.py diff --git a/gateway/__init__.py b/gateway/__init__.py index 173a052b57..7bf895717c 100644 --- a/gateway/__init__.py +++ b/gateway/__init__.py @@ -1,7 +1,7 @@ """Hermes Gateway - multi-platform messaging integration (sessions, context injection, delivery routing, platform-specific toolsets).""" -from .config import GatewayConfig, PlatformConfig, HomeChannel, SessionResetPolicy, load_gateway_config +from .config import GatewayConfig, PlatformConfig, HomeChannel, load_gateway_config from .session import ( SessionContext, SessionStore, @@ -11,6 +11,6 @@ from .delivery import DeliveryRouter, DeliveryTarget __all__ = [ "GatewayConfig", "PlatformConfig", "HomeChannel", "load_gateway_config", - "SessionContext", "SessionStore", "SessionResetPolicy", "build_session_context_prompt", + "SessionContext", "SessionStore", "build_session_context_prompt", "DeliveryRouter", "DeliveryTarget", ] diff --git a/gateway/config.py b/gateway/config.py index f75324d67c..de05751a4d 100644 --- a/gateway/config.py +++ b/gateway/config.py @@ -322,16 +322,15 @@ def persist_home_channel(home: HomeChannel, *, enabled_if_new: bool = False) -> @dataclass class SessionResetPolicy: - """When sessions reset: "daily" (at ``at_hour``), "idle" (after ``idle_minutes``), - "both" (whichever first), "none" (default: only compression manages context).""" + """Inert legacy value type retained solely for the scheduled plugin-compat window. + + Gateway configuration and session lifecycle do not consume this datatype. + """ mode: str = "none" at_hour: int = 4 # 0-23, local time idle_minutes: int = 1440 notify: bool = True # Notify the user when auto-reset occurs notify_exclude_platforms: tuple = ("api_server", "webhook") - # A background process this old no longer blocks reset (not killed, only ignored by the guard). - # A forgotten preview server should not keep a session alive forever (#29177). Raise this if you run - # legitimate multi-day jobs whose liveness should pin the conversation open. bg_process_max_age_hours: int = 24 def to_dict(self) -> Dict[str, Any]: @@ -540,9 +539,6 @@ _TOPLEVEL_BOOL_DEFAULTS = { class GatewayConfig: """Main gateway configuration: platform connections, session policies, delivery settings.""" platforms: Dict[Platform, PlatformConfig] = field(default_factory=dict) - default_reset_policy: SessionResetPolicy = field(default_factory=SessionResetPolicy) - reset_by_type: Dict[str, SessionResetPolicy] = field(default_factory=dict) - reset_by_platform: Dict[Platform, SessionResetPolicy] = field(default_factory=dict) reset_triggers: List[str] = field(default_factory=lambda: ["/new", "/reset"]) quick_commands: Dict[str, Any] = field(default_factory=dict) # slash commands that bypass the agent loop sessions_dir: Path = field(default_factory=lambda: get_hermes_home() / "sessions") @@ -650,20 +646,9 @@ class GatewayConfig: def get_home_channel(self, platform: Platform) -> Optional[HomeChannel]: return self.platforms[platform].home_channel if self.platforms.get(platform) else None - def get_reset_policy(self, platform: Optional[Platform] = None, session_type: Optional[str] = None) -> SessionResetPolicy: - """Priority: platform override > type override > default.""" - if platform and platform in self.reset_by_platform: - return self.reset_by_platform[platform] - if session_type and session_type in self.reset_by_type: - return self.reset_by_type[session_type] - return self.default_reset_policy - def to_dict(self) -> Dict[str, Any]: return { "platforms": {p.value: c.to_dict() for p, c in self.platforms.items()}, - "default_reset_policy": self.default_reset_policy.to_dict(), - "reset_by_type": {k: v.to_dict() for k, v in self.reset_by_type.items()}, - "reset_by_platform": {p.value: v.to_dict() for p, v in self.reset_by_platform.items()}, "reset_triggers": self.reset_triggers, "quick_commands": self.quick_commands, "sessions_dir": str(self.sessions_dir), @@ -739,14 +724,6 @@ class GatewayConfig: return cls( platforms=by_platform("platforms", PlatformConfig.from_dict, dicts_only=True), - default_reset_policy=SessionResetPolicy.from_dict(data["default_reset_policy"]) - if "default_reset_policy" in data - else SessionResetPolicy(), - reset_by_type={ - type_name: SessionResetPolicy.from_dict(policy_data) - for type_name, policy_data in _coerce_dict(data.get("reset_by_type", {})).items() - }, - reset_by_platform=by_platform("reset_by_platform", SessionResetPolicy.from_dict), reset_triggers=data.get("reset_triggers", ["/new", "/reset"]), quick_commands=_coerce_dict(data.get("quick_commands", {})), sessions_dir=Path(data["sessions_dir"]) if "sessions_dir" in data else get_hermes_home() / "sessions", @@ -820,16 +797,6 @@ def load_gateway_config() -> GatewayConfig: def _validate_gateway_config(config: "GatewayConfig") -> None: """Validate and sanitize a loaded GatewayConfig in place (after all sources are merged).""" - policy = config.default_reset_policy - - if not (0 <= policy.at_hour <= 23): - logger.warning("Invalid at_hour=%s (must be 0-23). Using default 4.", policy.at_hour) - policy.at_hour = 4 - - if policy.idle_minutes is None or policy.idle_minutes <= 0: - logger.warning("Invalid idle_minutes=%s (must be positive). Using default 1440.", policy.idle_minutes) - policy.idle_minutes = 1440 - try: # Reject known-weak placeholder tokens. Ported from openclaw/openclaw#64586: users who copy # .env.example without changing placeholder values get a clear startup error instead of a confusing diff --git a/gateway/config_env.py b/gateway/config_env.py index a25c4b8fea..a1e80ec0be 100644 --- a/gateway/config_env.py +++ b/gateway/config_env.py @@ -341,13 +341,6 @@ def _qq_home(config: GatewayConfig, qq_config: PlatformConfig) -> None: ) -def _session_settings(config: GatewayConfig) -> None: - for env, attr in (("SESSION_IDLE_MINUTES", "idle_minutes"), ("SESSION_RESET_HOUR", "at_hour")): - if raw := getenv(env): - with contextlib.suppress(ValueError): - setattr(config.default_reset_policy, attr, int(raw)) - - def _plugin_probe_seed(entry) -> Optional[dict]: """``env_enablement_fn()`` result as a non-empty dict, else None.""" if entry.env_enablement_fn is None: @@ -623,7 +616,7 @@ _ENV_STEPS: tuple = ( ), home="YUANBAO_HOME_CHANNEL", ), - _session_settings, + _enable_plugin_platforms_from_env, _relay, _scrub_explicit_markers, diff --git a/gateway/config_loader.py b/gateway/config_loader.py index 52167d67bd..aa279b3608 100644 --- a/gateway/config_loader.py +++ b/gateway/config_loader.py @@ -68,7 +68,6 @@ def _presence(*keys: str) -> tuple: # (yaml key, gw_data key, mode, accept(value) -> bool, transform(value)) _TOPLEVEL_BRIDGE: tuple = ( - ("session_reset", "default_reset_policy", "presence", lambda v: bool(v) and isinstance(v, dict), None), ("quick_commands", "quick_commands", "none", _quick_commands_ok, None), ("stt", "stt", "presence", lambda v: isinstance(v, dict), None), *_presence("stt_echo_transcripts", "group_sessions_per_user", "thread_sessions_per_user"), diff --git a/gateway/run.py b/gateway/run.py index d6bfc83ec1..7557085b25 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -41,7 +41,7 @@ from agent.turn_context import compression_made_progress from hermes_cli.config import _is_ssh_remote_tilde_cwd, cfg_get from hermes_cli.fallback_config import get_fallback_chain -# Per-session AIAgent cache bounds (agents are heavy); see _enforce_agent_cache_cap/_session_expiry_watcher. +# Per-session AIAgent cache bounds (agents are heavy); see _enforce_agent_cache_cap/_session_housekeeping_watcher. _AGENT_CACHE_MAX_SIZE = 128 _AGENT_CACHE_IDLE_TTL_SECS = 3600.0 # evict agents idle for >1h _PLATFORM_CONNECT_TIMEOUT_SECS_DEFAULT = 30.0 @@ -3395,16 +3395,11 @@ class GatewayRunner( def _init_session_store(self) -> None: """Build the SessionStore (with process-registry reset guard), its async facade and the router.""" - # Reset guard: a background process older than session_reset.bg_process_max_age_hours (24h - # default) is stale and no longer blocks idle/daily reset (NOT killed, only ignored). from tools.process_registry import process_registry - _bg_max_age_hours = getattr(self.config.default_reset_policy, "bg_process_max_age_hours", 24) - _bg_max_age_seconds = ( - _bg_max_age_hours * 3600 if _bg_max_age_hours and _bg_max_age_hours > 0 else None) self.session_store = SessionStore( self.config.sessions_dir, self.config, has_active_processes_fn=lambda key: process_registry.has_active_for_session( - key, max_active_age=_bg_max_age_seconds)) + key)) # Loop-side boundary: sync helpers use ``session_store`` directly; async handlers await this facade. self._async_session_store = AsyncSessionStore(self.session_store) self.delivery_router = DeliveryRouter(self.config) diff --git a/gateway/run_agent_cache.py b/gateway/run_agent_cache.py index cac5714863..7ffa809afe 100644 --- a/gateway/run_agent_cache.py +++ b/gateway/run_agent_cache.py @@ -612,36 +612,17 @@ class GatewayAgentCacheMixin: with suppress(Exception): target(*args) - def _finalizable_unexpired_session_entry(self, key: str): - """Session-store entry for ``key`` when the expiry watcher will still finalize it; None when - missing, not finalizable (``mode == "none"``) or already expired (the watcher handles those).""" - _store = getattr(self, "session_store", None) - if _store is None: - return None - try: - _store._ensure_loaded() - entry = _store._entries.get(key) - except Exception: - return None - ok = entry is not None and _store.is_session_finalizable(entry) and not _store._is_session_expired(entry) - return entry if ok else None - def _commit_memory_before_soft_evict(self, agent: Any, key: str) -> None: - """Fire on_session_end extraction before soft-evicting a live agent: the expiry watcher only - finalizes what it finds in ``_agent_cache``, so an LRU soft-evict first would hide the - transcript from memory providers. Commit via ``commit_memory_session`` (no teardown), only for - finalizable, not-yet-expired sessions. Best-effort.""" + """Commit the live transcript to memory providers before resource-only eviction.""" # No external memory provider (``_memory_manager`` None) — nothing to commit. if agent is None or not hasattr(agent, "commit_memory_session") or getattr(agent, "_memory_manager", None) is None: return try: - if self._finalizable_unexpired_session_entry(key) is None: - return messages = getattr(agent, "_session_messages", None) agent.commit_memory_session(messages if isinstance(messages, list) else None) logger.debug( "Committed on_session_end extraction before soft-evicting " - "finalizable session=%s (cache pressure, pre-expiry)", key, + "session=%s (resource eviction)", key, ) except Exception as _e: logger.debug("Pre-evict memory commit failed for %s: %s", key, _e) @@ -846,17 +827,10 @@ class GatewayAgentCacheMixin: last_activity = getattr(agent, "_last_activity_ts", None) if last_activity is None or (now - last_activity) <= idle_ttl: continue - # Not yet expired in the store (daily-reset fires hours after the last message): keep - # the agent so the expiry watcher can call on_session_end() with the live transcript. - # Only defer when the watcher will EVER finalize it — for mode == "none" deferring pins - # the agent for the gateway's lifetime (the leak this sweep relieves); those soft-evict - # WITHOUT on_session_end, correctly. - if self._finalizable_unexpired_session_entry(key) is not None: - continue to_evict.append((key, agent)) for key, _ in to_evict: _cache.pop(key, None) for key, agent in to_evict: logger.info("Agent cache idle-TTL evict: session=%s (idle=%.0fs)", key, now - getattr(agent, "_last_activity_ts", now)) - self._spawn_release_thread(self._release_evicted_agent_soft, (agent,), f"agent-cache-idle-{key[:24]}", inline_fallback=False) + self._spawn_release_thread(self._commit_then_release_soft, (agent, key), f"agent-cache-idle-{key[:24]}", inline_fallback=False) return len(to_evict) diff --git a/gateway/run_startup.py b/gateway/run_startup.py index 24d187b84b..72fe624024 100644 --- a/gateway/run_startup.py +++ b/gateway/run_startup.py @@ -1267,7 +1267,7 @@ class GatewayStartupMixin: # Long-lived supervised watchers spawned at the end of start(), in order; supervised name = method # name minus the leading underscore. _PRE_RECONNECT_WATCHERS = ( - "_session_expiry_watcher", "_model_catalog_refresh_watcher", "_session_stall_watcher", + "_session_housekeeping_watcher", "_model_catalog_refresh_watcher", "_session_stall_watcher", "_kanban_notifier_watcher", "_kanban_dispatcher_watcher", ) _POST_RECONNECT_WATCHERS = ("_handoff_watcher", "_async_delegation_watcher", "_loop_wakeup_watcher") diff --git a/gateway/run_turn.py b/gateway/run_turn.py index dc6e59997c..88cad72fc9 100644 --- a/gateway/run_turn.py +++ b/gateway/run_turn.py @@ -387,9 +387,9 @@ class GatewayTurnMixin: async def _hmwa_deliver_auto_reset_notice(self, session_entry, source, turn_sidecar_notes): """Stage the auto-reset sidecar note for the agent and notify the user (policy-gated).""" - from gateway.run import _AUTO_RESET_CONTEXT_NOTES, _auto_reset_reason_text - reset_reason = getattr(session_entry, 'auto_reset_reason', None) or 'idle' - context_note = _AUTO_RESET_CONTEXT_NOTES.get(reset_reason, _AUTO_RESET_CONTEXT_NOTES["idle"]) + from gateway.run import _AUTO_RESET_CONTEXT_NOTES + reset_reason = getattr(session_entry, 'auto_reset_reason', None) or 'suspended' + context_note = _AUTO_RESET_CONTEXT_NOTES.get(reset_reason, _AUTO_RESET_CONTEXT_NOTES["suspended"]) # Long-lived channels: point the agent at the prior same-channel session for session_search. try: # Returns None (appends nothing) for other platforms or when there's no prior activity to @@ -402,34 +402,13 @@ class GatewayTurnMixin: turn_sidecar_notes.append(context_note) try: - policy = self.session_store.config.get_reset_policy( - platform=source.platform, session_type=getattr(source, 'chat_type', 'dm'), - ) - # Check pairing store. A pairing entry is a first-class authorization grant, created only by a - # trusted operator approving a pairing code (hermes gateway pairing approve / the authenticated - # dashboard) — an inbound sender can never reach approve_code, so this is not an - # attacker-controlled path. Honored as a UNION with the allowlist: a paired user is authorized - # regardless of the allowlist, and when an allowlist IS configured, operator approval also - # writes the user into that allowlist (see PairingStore._approve_user), keeping a single - # operator-visible source of truth. (#23778: the original bypass was the inbound - # message/approval-button gate, not this gate; that gate is fixed separately.) In multiplex - # gateways, route to the per-profile PairingStore so each profile's whitelist is isolated; falls - # back to the global store when the source has no profile or the profile isn't registered. - platform_name = source.platform.value if source.platform else "" - # Suspended / restart-recovery-expired sessions always notify (the user must learn they - # can /resume); idle/daily resets respect policy.notify + excluded platforms + activity. - should_notify = reset_reason in {"suspended", "resume_pending_expired"} or ( - policy.notify - and getattr(session_entry, 'reset_had_activity', False) - and platform_name not in policy.notify_exclude_platforms - ) + should_notify = reset_reason == "suspended" adapter = self._adapter_for_source(source) if should_notify else None if adapter: notice = ( - f"◐ Session automatically reset ({_auto_reset_reason_text(reset_reason, policy)}). " + "◐ Session reset after being stopped. " f"Conversation history cleared.\n" f"Use /resume to browse and restore a previous session.\n" - f"Adjust reset timing in config.yaml under session_reset." ) with suppress(Exception): session_info = await asyncio.to_thread(self._reset_notice_session_info, source) diff --git a/gateway/run_watchers.py b/gateway/run_watchers.py index ec7493eea4..a84feada51 100644 --- a/gateway/run_watchers.py +++ b/gateway/run_watchers.py @@ -23,15 +23,9 @@ from gateway.session_stall import ( # Log-record parity with the origin module. logger = logging.getLogger("gateway.run") -_MAX_FINALIZE_RETRIES = 3 _SESSION_STORE_PRUNE_INTERVAL = 3600.0 # once per hour -def _platform_of_key(key: str, default: str = "") -> str: - """Session keys look like ``agent:main:telegram:dm:12345`` — platform is field [2].""" - return key.split(":")[2] if key.count(":") > 1 else default - - async def _interruptible_sleep(runner, seconds: int) -> None: """Sleep in 1s increments so the watcher stops quickly when ``runner._running`` flips.""" for _ in range(seconds): @@ -43,88 +37,17 @@ async def _interruptible_sleep(runner, seconds: int) -> None: class GatewaySessionWatchersMixin: """Session expiry / stall / catalog-refresh watcher loops for GatewayRunner.""" - async def _session_expiry_watcher(self, interval: int = 300): - """Finalize expired sessions, then run the cache/store sweeps.""" - await asyncio.sleep(60) # initial delay — let the gateway fully start - finalize_failures: dict[str, int] = {} # session_id -> consecutive failure count + async def _session_housekeeping_watcher(self, interval: int = 300): + """Reclaim resources without ending durable conversations.""" + await asyncio.sleep(60) while self._running: try: - store = self.async_session_store - await store._ensure_loaded() - expired = [ - (key, entry) for key, entry in list(self.session_store._entries.items()) - if not entry.expiry_finalized and await store._is_session_expired(entry) - ] - if expired: - await self._finalize_expired_sessions(expired, finalize_failures) - await self._expiry_housekeeping() + await self._session_housekeeping() except Exception as e: - logger.debug("Session expiry watcher error: %s", e) + logger.debug("Session housekeeping error: %s", e) await _interruptible_sleep(self, interval) - async def _finalize_expired_sessions(self, expired: list, failures: dict[str, int]) -> None: - """Finalize each entry; after ``_MAX_FINALIZE_RETRIES`` consecutive failures mark it - finalized anyway (without clearing the model override) to stop an infinite retry loop.""" - platforms = Counter(_platform_of_key(k, "unknown") for k, _ in expired) - logger.info( - "Session expiry: %d sessions to finalize (%s)", - len(expired), ", ".join(f"{p}:{c}" for p, c in sorted(platforms.items())), - ) - for key, entry in expired: - sid = entry.session_id - try: - await self._finalize_expired_session(key, entry) - except Exception as e: - count = failures[sid] = failures.get(sid, 0) + 1 - if count < _MAX_FINALIZE_RETRIES: - logger.debug("Session finalize failed (%d/%d) for %s: %s", - count, _MAX_FINALIZE_RETRIES, sid, e) - continue - logger.warning( - "Session finalize gave up after %d attempts for %s: %s. " - "Marking as finalized to prevent infinite retry loop.", count, sid, e, - ) - store = self.async_session_store - await store.set_expiry_finalized(entry, clear_model_override=False) - failures.pop(sid, None) - done = sum(1 for _, e in expired if e.expiry_finalized) - if failed := len(expired) - done: - logger.info("Session expiry done: %d finalized, %d pending retry", done, failed) - else: - logger.info("Session expiry done: %d finalized", done) - - async def _finalize_expired_session(self, key: str, entry) -> None: - """Run finalize hooks, tear down the cached agent, clear conversation scope, persist.""" - from gateway.run import _AGENT_PENDING_SENTINEL - # Off-loop + bounded: plugin finalize hooks can block arbitrarily; this is the event loop. - with contextlib.suppress(Exception): - await self._finalize_session_off_loop( - session_id=entry.session_id, platform=_platform_of_key(key), - reason="session_expired", - ) - # Idle agents live in _agent_cache (not _running_agents); fall back to the running turn's - # agent in case the session is still mid-turn when the expiry fires. - agent, cache_lock = None, getattr(self, "_agent_cache_lock", None) # tests may lack it - if cache_lock is not None: - with cache_lock: - cached = self._agent_cache.get(key) - agent = cached[0] if isinstance(cached, tuple) else cached if cached else None - if agent is None: - state = self._peek_session_state(key) - agent = state.turn.agent if state else None - if agent and agent is not _AGENT_PENDING_SENTINEL: - await self._cleanup_agent_resources_off_loop(agent, context="session expiry") - # Evict so the AIAgent (LLM clients, tool schemas, memory refs) can be GC'd, then drop - # every conversation-scoped dict AND boundary security state — only finalize, /new, /reset - # may (idle-cache eviction must NOT: a resumed turn rebuilds from those overrides). The - # persisted flag also drops the /model override: finalization is a conversation boundary. - self._evict_cached_agent(key) - self._clear_conversation_scope(key, reason="expiry_finalized") - # See #9006. - await self.async_session_store.set_expiry_finalized(entry) - logger.debug("Session expiry finalized for %s", entry.session_id) - - async def _expiry_housekeeping(self) -> None: + async def _session_housekeeping(self) -> None: """Idle/pressure agent-cache sweeps plus the hourly SessionStore prune.""" # Sweep idle agents regardless of reset policy: long/"never" windows would pin memory. try: diff --git a/gateway/session.py b/gateway/session.py index e07f49a11b..3a1fc4340f 100644 --- a/gateway/session.py +++ b/gateway/session.py @@ -901,7 +901,7 @@ class SessionStore( sid = observed.session_id checks = _RouteChecks( sid, self._compression_tip_for_session_id(sid), self._is_session_ended_in_db(sid), - self._route_reset_reason(observed, source, now), + self._route_reset_reason(observed), ) # Phase 2 (lock): apply the decisions to _entries. decision = self._apply_route_checks(session_key, checks, force_new, touch_activity, now) @@ -981,10 +981,6 @@ class SessionStore( recovered = self._query_recoverable_session(session_key=session_key, source=source, now=now) if recovered is None: return - reset_reason = self._should_reset(recovered, source) - if reset_reason: - decision.schedule_reset(reset_reason, recovered, recovered.reset_had_activity) - return self._reopen_session_row(session_key, recovered.session_id) with self._lock: decision.entry = self._entries.setdefault(session_key, recovered) diff --git a/gateway/session_lifecycle.py b/gateway/session_lifecycle.py index 78fb67e3ae..1c037af626 100644 --- a/gateway/session_lifecycle.py +++ b/gateway/session_lifecycle.py @@ -1,6 +1,4 @@ -"""SessionStore reset/expiry policy and crash-recovery markers (idle/daily reset, expiry -finalization, active-turn tokens, resume_pending, suspension, pruning) plus the shared clock/id -helpers. Mixin split out of ``gateway/session.py``; bound onto ``SessionStore`` via the MRO.""" +"""SessionStore explicit suspension, crash-recovery markers, pruning and shared clock/id helpers.""" from __future__ import annotations @@ -46,8 +44,7 @@ _AUTO_CONTINUE_FRESHNESS_SECS_DEFAULT = 60 * 60 def auto_continue_freshness_window() -> float: - """Auto-continue freshness window in seconds (one source of truth for the resume scheduler and - the routing-time zombie gate); env var, default when unset/malformed; non-positive disables.""" + """Resume-scheduler freshness window; stale automation never discards the transcript.""" raw = os.environ.get("HERMES_AUTO_CONTINUE_FRESHNESS") try: return float(raw) if raw else float(_AUTO_CONTINUE_FRESHNESS_SECS_DEFAULT) @@ -56,72 +53,7 @@ def auto_continue_freshness_window() -> float: class SessionLifecycleMixin: - """SessionStore reset/expiry policy and crash-recovery markers.""" - - def set_expiry_finalized( - self, entry: SessionEntry, *, clear_model_override: bool = True - ) -> None: - """Mark a session entry expiry-finalized in memory, sessions.json, AND state.db (single - write-path for the expiry watcher). ``clear_model_override=False`` = flag only.""" - with self._lock: - entry.expiry_finalized = True - if clear_model_override: - # Finalization is a conversation boundary: a later message must not rehydrate it. - entry.model_override = None - self._save() - # Background caller never entered ``_profile_runtime_scope``: resolve the store by key. - # The expiry watcher calls this from a background task that never entered - # ``_profile_runtime_scope``, so resolve the store from the key rather than from the ambient scope - # (#66887). - _db = self._db_for_key(entry.session_key) - if not _db: - return - setter = getattr(_db, "set_expiry_finalized", None) - if callable(setter): - try: - setter(entry.session_id, True) - except Exception as exc: - logger.debug("Session DB expiry_finalized write failed for %s: %s", entry.session_id, exc) - try: - # Without a durable ``session_reset`` end_reason, later agent cleanup ends the row as - # ``agent_close``, which stale-route recovery treats as resumable. - _db.promote_to_session_reset(entry.session_id) - except Exception as exc: - logger.debug("Session DB promote_to_session_reset failed for %s: %s", entry.session_id, exc) - - @staticmethod - def _policy_reset_reason(policy, updated_at: datetime) -> Optional[str]: - """Return "idle"/"daily" when *updated_at* is overdue under *policy*, else None.""" - if policy.mode == "none": - return None - now = _now() - if policy.mode in {"idle", "both"} and now > updated_at + timedelta(minutes=policy.idle_minutes): - return "idle" - if policy.mode in {"daily", "both"}: - today_reset = now.replace(hour=policy.at_hour, minute=0, second=0, microsecond=0) - if now.hour < policy.at_hour: - today_reset -= timedelta(days=1) - if updated_at < today_reset: - return "daily" - return None - - def _is_session_expired(self, entry: SessionEntry) -> bool: - """Whether the reset policy has expired *entry* (expiry watcher); sessions with active - background processes never expire.""" - if self._has_active_processes_safe(entry.session_key, context="expiry"): - logger.debug("Session %s not expired — active background processes", entry.session_key) - return False - policy = self.config.get_reset_policy(platform=entry.platform, session_type=entry.chat_type) - return self._policy_reset_reason(policy, entry.updated_at) is not None - - def is_session_finalizable(self, entry: SessionEntry) -> bool: - """True if the expiry watcher will *ever* finalize this session; ``mode == "none"`` never - expires, so the agent-cache sweep must reap its agent itself. Policy errors -> False.""" - try: - policy = self.config.get_reset_policy(platform=entry.platform, session_type=entry.chat_type) - return policy.mode != "none" - except Exception: - return False + """SessionStore explicit boundaries and crash-recovery markers.""" def _is_session_ended_in_db(self, session_id: str) -> bool: """True iff state.db has this session with a non-null end_reason (same staleness test as @@ -146,37 +78,9 @@ class SessionLifecycleMixin: return False return bool(row is not None and row.get("end_reason") is not None) - def _should_reset(self, entry: SessionEntry, source: SessionSource) -> Optional[str]: - """Reset reason ("idle"/"daily") if policy says reset, else None; sessions with active - background processes are never reset.""" - session_key = self._generate_session_key(source) - if self._has_active_processes_safe(session_key, context="reset"): - logger.debug("Session reset skipped for %s — active background processes", session_key) - return None - policy = self.config.get_reset_policy(platform=source.platform, session_type=source.chat_type) - return self._policy_reset_reason(policy, entry.updated_at) - - def _route_reset_reason( - self, entry: SessionEntry, source: SessionSource, now: datetime - ) -> Optional[str]: - """Reset decision for an existing route (no lock; DB/config I/O). ``suspended`` always - resets; otherwise the reset policy decides, and a still-pending resume marker is also - freshness-gated — but ``session_reset.mode: none`` (user opted out of ALL automatic - resets) makes an expired marker fall through to a normal resume, never a silent fresh - session.""" - if entry.suspended: - return "suspended" - reason = self._should_reset(entry, source) - if reason or not entry.resume_pending: - return reason - policy = self.config.get_reset_policy(platform=source.platform, session_type=source.chat_type) - if policy.mode == "none": - return None - window = auto_continue_freshness_window() - ref_time = entry.last_resume_marked_at or entry.updated_at - if window > 0 and (now - ref_time).total_seconds() > window: - return "resume_pending_expired" - return None + def _route_reset_reason(self, entry: SessionEntry) -> Optional[str]: + """Only explicit suspension replaces a routed conversation; time never does.""" + return "suspended" if entry.suspended else None def _update_entry(self, session_key: str, mutate) -> bool: """Apply ``mutate(entry)`` under ``_lock`` and full-save; False when the entry is missing diff --git a/gateway/session_persistence.py b/gateway/session_persistence.py index 95c0ef5757..f9b8ad3dc1 100644 --- a/gateway/session_persistence.py +++ b/gateway/session_persistence.py @@ -146,7 +146,7 @@ class SessionPersistenceMixin: into the ROOT store until the stale-route self-heal drops a live conversation. Background work runs unscoped while operating on every profile's keys out of the single process-wide - ``_entries`` dict — ``_session_expiry_watcher`` is the clearest case — so it reads and writes the + ``_entries`` dict — ``_session_housekeeping_watcher`` is the clearest case — so it reads and writes the ROOT store for rows that actually live under ``profiles//state.db``. The two writers then drift apart on the same logical session until the routing index disagrees with the row and the #54878 self-heal drops a live conversation (#66887). diff --git a/gateway/session_recovery.py b/gateway/session_recovery.py index ffe341c8fb..4b7219428b 100644 --- a/gateway/session_recovery.py +++ b/gateway/session_recovery.py @@ -181,9 +181,7 @@ class SessionRecoveryMixin: def _recover_session_from_db( self, *, session_key: str, source: SessionSource, now: datetime, raise_on_lookup_error: bool = False) -> Optional[SessionEntry]: - """Rebuild a missing session-key mapping from durable state.db data. ``None`` when no row is - recoverable, or when the recovered session is already overdue under the reset policy — the - row is then durably promoted to a reset boundary instead of resurrected.""" + """Rebuild a missing session-key mapping from a recoverable durable row.""" entry, migrated_legacy = self._query_recoverable_row( # The legacy (pre-workspace) Slack key fallback happens INSIDE _query_recoverable_session # (#20583/#66398 design): it performs the exact-key legacy lookup, claims the key once per @@ -192,15 +190,6 @@ class SessionRecoveryMixin: raise_on_lookup_error=raise_on_lookup_error) if entry is None: return None - reset_reason = self._should_reset(entry, source) - if reset_reason: - self._promote_session_reset( - session_key, entry.session_id, reset_reason, - log=lambda exc: logger.debug( - "Gateway recovered-session reset promotion failed for %s: %s", session_key, exc, - ), - ) - return None self._reopen_session_row(session_key, entry.session_id) if migrated_legacy: self._record_gateway_session_peer( diff --git a/hermes_cli/config.py b/hermes_cli/config.py index 9dfc5e0b9d..b70ff27022 100644 --- a/hermes_cli/config.py +++ b/hermes_cli/config.py @@ -1062,7 +1062,7 @@ _EXTRA_KNOWN_ROOT_KEYS = { "known_builtin_toolsets", # ditto — builtin toolsets a platform's checklist has offered "tool_gateway_declined_tools", # per-tool Tool Gateway offer declines # Top-level forms read/bridged by gateway/config.py: - "session_reset", "group_sessions_per_user", "thread_sessions_per_user", + "group_sessions_per_user", "thread_sessions_per_user", "stt_echo_transcripts", "reset_triggers", "always_log_local", "filter_silence_narration", "multiplex_profiles", "profile_routes", "platforms", "require_mention", "unauthorized_dm_behavior", "signal", diff --git a/hermes_cli/setup.py b/hermes_cli/setup.py index 0feb31f3b5..39ab228613 100644 --- a/hermes_cli/setup.py +++ b/hermes_cli/setup.py @@ -403,12 +403,9 @@ def _apply_default_agent_settings(config: dict): config.setdefault("display", {})["tool_progress"] = "all" config.setdefault("compression", {})["enabled"] = True config["compression"]["threshold"] = 0.50 - # Never auto-reset (the gateway default); written explicitly so it is visible in config.yaml. - config.setdefault("session_reset", {})["mode"] = "none" save_config(config) print_success("Applied recommended defaults:") _info(" Max iterations: 150", " Tool progress: all", " Compression threshold: 0.50", - " Session reset: never (use /reset or compression)", " Run `hermes setup agent` later to customize.") @@ -435,24 +432,6 @@ _TOOL_PROGRESS_HELP = ( " verbose — Full args, results, and debug logs", " log — Silent in chat; write every tool call to ~/.hermes/logs/tool_calls.log (gateway only)", ) -_SESSION_RESET_HELP = ( - "Messaging sessions (Telegram, Discord, etc.) accumulate context over time.", - "Each message adds to the conversation history, which means growing API costs.", "", - "To manage this, sessions can automatically reset after a period of inactivity", - "or at a fixed time each day. When a reset happens, the agent saves important", - "things to its persistent memory first — but the conversation context is cleared.", "", - "You can also manually reset anytime by typing /reset in chat.", "", -) -_SESSION_RESET_CHOICES = [ - "Inactivity + daily reset (reset whichever comes first)", - "Inactivity only (reset after N minutes of no messages)", - "Daily only (reset at a fixed hour each day)", - "Never auto-reset (recommended - context lives until /reset or context compression)", - "Keep current settings", -] -_SESSION_RESET_MODES = ("both", "idle", "daily", "none") # index 4 = keep current - - def setup_agent_settings(config: dict): """Configure agent behavior: iterations, progress display, compression, session reset.""" print_header("Agent Settings") @@ -497,39 +476,9 @@ def setup_agent_settings(config: dict): config["compression"]["threshold"] = threshold print_success(f"Context compression threshold set to {config['compression'].get('threshold', 0.50)}") - # ── Session Reset Policy ── - print_header("Session Reset Policy") - _info(*_SESSION_RESET_HELP) - _prompt_session_reset(config.setdefault("session_reset", {})) save_config(config) -def _prompt_session_reset(reset_cfg: dict) -> None: - """Pick the session reset mode and its idle/daily parameters in place.""" - current_mode = reset_cfg.get("mode", "none") - current_idle, current_hour = reset_cfg.get("idle_minutes", 1440), reset_cfg.get("at_hour", 4) - default_reset = _SESSION_RESET_MODES.index(current_mode) if current_mode in _SESSION_RESET_MODES else 3 - reset_idx = prompt_choice("Session reset mode:", _SESSION_RESET_CHOICES, default_reset) - mode = _SESSION_RESET_MODES[reset_idx] if 0 <= reset_idx < len(_SESSION_RESET_MODES) else None - if mode is None: # keep current settings - return - reset_cfg["mode"] = mode - if mode in ("both", "idle"): - _prompt_int_setting(reset_cfg, "idle_minutes", " Inactivity timeout (minutes)", current_idle, lambda v: v > 0) - if mode in ("both", "daily"): - _prompt_int_setting(reset_cfg, "at_hour", " Daily reset hour (0-23, local time)", current_hour, lambda v: 0 <= v <= 23) - idle_now, hour_now = reset_cfg.get("idle_minutes", 1440), reset_cfg.get("at_hour", 4) - if mode == "none": - print_info("Sessions will never auto-reset. Context is managed only by compression.") - print_warning("Long conversations will grow in cost. Use /reset manually when needed.") - else: - print_success({ - "both": f"Sessions reset after {idle_now} min idle or daily at {hour_now}:00", - "idle": f"Sessions reset after {idle_now} min of inactivity", - "daily": f"Sessions reset daily at {hour_now}:00", - }[mode]) - - # ── Section 5: Tool Configuration (delegates to unified tools_config.py) ── diff --git a/hermes_cli/setup_quick.py b/hermes_cli/setup_quick.py index 9f282af1a7..5f31223aa1 100644 --- a/hermes_cli/setup_quick.py +++ b/hermes_cli/setup_quick.py @@ -194,7 +194,6 @@ def _blank_slate_minimize_config(config: dict): mem["user_profile_enabled"] = False config.setdefault("checkpoints", {})["enabled"] = False config.setdefault("smart_model_routing", {})["enabled"] = False - config.setdefault("session_reset", {})["mode"] = "none" config.setdefault("display", {})["tool_progress"] = "all" diff --git a/tests/gateway/test_session_time_persistence.py b/tests/gateway/test_session_time_persistence.py new file mode 100644 index 0000000000..b96b983e1b --- /dev/null +++ b/tests/gateway/test_session_time_persistence.py @@ -0,0 +1,41 @@ +"""Elapsed time cannot replace a durable conversation; explicit boundaries still can.""" +from datetime import datetime, timedelta + +import pytest + +from gateway.config import GatewayConfig, Platform +from gateway.session import SessionSource, SessionStore + + +@pytest.mark.parametrize("mode", ["idle", "daily", "both", "none"]) +def test_old_reset_config_cannot_rotate_durable_conversation(tmp_path, mode): + config = GatewayConfig.from_dict({ + "default_reset_policy": {"mode": mode, "idle_minutes": 1}, + "reset_by_type": {"dm": {"mode": mode, "idle_minutes": 1}}, + "reset_by_platform": {"telegram": {"mode": mode, "idle_minutes": 1}}, + }) + store = SessionStore(tmp_path / "sessions", config) + source = SessionSource(platform=Platform.TELEGRAM, chat_id="time-invariant", user_id="test") + old = store.get_or_create_session(source) + messages = [{"role": "user", "content": "keep my conversation"}, + {"role": "assistant", "content": "including after a restart"}] + for message in messages: + store.append_to_transcript(old.session_id, message) + old.updated_at = datetime.now() - timedelta(days=3) + old.resume_pending = True + old.last_resume_marked_at = old.updated_at + store._save() + routed = store.get_or_create_session(source) + assert routed.session_id == old.session_id + assert store.load_transcript(routed.session_id) == messages + assert store._db.get_session(old.session_id)["end_reason"] is None + # Missing routing indexes must recover the same durable transcript too. + store._entries.clear() + store._save() + recovered = store.get_or_create_session(source) + assert recovered.session_id == old.session_id + explicit = store.reset_session(recovered.session_key) + assert explicit.session_id != old.session_id + assert store._db.get_session(old.session_id)["end_reason"] == "session_reset" + assert store.load_transcript(old.session_id) == messages + store._db.close()