fix(gateway): stop time-triggered conversation rotation

This commit is contained in:
Teknium
2026-09-07 03:22:27 -07:00
parent 7798241eab
commit 1d5d059410
17 changed files with 75 additions and 367 deletions
+2 -2
View File
@@ -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",
]
+4 -37
View File
@@ -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
+1 -8
View File
@@ -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,
-1
View File
@@ -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"),
+2 -7
View File
@@ -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)
+3 -29
View File
@@ -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)
+1 -1
View File
@@ -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")
+5 -26
View File
@@ -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)
+6 -83
View File
@@ -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:
+1 -5
View File
@@ -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)
+6 -102
View File
@@ -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
+1 -1
View File
@@ -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/<name>/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).
+1 -12
View File
@@ -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(
+1 -1
View File
@@ -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",
-51
View File
@@ -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) ──
-1
View File
@@ -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"
@@ -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()