refactor(gateway): session — rewrap comment/docstring paragraphs to 100 cols
This commit is contained in:
+29
-41
@@ -40,8 +40,7 @@ from gateway.session_transcript import SessionTranscriptMixin, TranscriptReadErr
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# PII redaction helpers
|
||||
# --------------------------------------------------------------------------- PII redaction helpers
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def _hash_id(value: str) -> str:
|
||||
@@ -98,10 +97,9 @@ class SessionSource:
|
||||
user_id_alt: Optional[str] = None # platform-specific stable alt ID (Signal UUID, Feishu union_id)
|
||||
chat_id_alt: Optional[str] = None # Signal group internal ID
|
||||
is_bot: bool = False # message author is a bot/webhook (Discord)
|
||||
# Platform-neutral SCOPE discriminator (Discord guild / Slack workspace /
|
||||
# Matrix server) driving server/workspace isolation. ``scope_id`` is
|
||||
# canonical; ``guild_id`` is a deprecated alias kept while the cross-repo
|
||||
# dual-read/dual-write overlap lasts (both written, scope_id wins on read).
|
||||
# Platform-neutral SCOPE discriminator (Discord guild / Slack workspace / Matrix server) driving
|
||||
# server/workspace isolation. ``scope_id`` is canonical; ``guild_id`` is a deprecated alias kept
|
||||
# while the cross-repo dual-read/dual-write overlap lasts (both written, scope_id wins on read).
|
||||
scope_id: Optional[str] = None
|
||||
guild_id: Optional[str] = None
|
||||
parent_chat_id: Optional[str] = None # parent channel when chat_id is a thread
|
||||
@@ -121,11 +119,10 @@ class SessionSource:
|
||||
# that WILL be delivered into a new thread whose id == this message id, so
|
||||
# the initiating message and later in-thread follow-ups share ONE session.
|
||||
prospective_thread_id: Optional[str] = None
|
||||
# Wire-INVISIBLE trust signal: delivered over the per-instance authenticated
|
||||
# relay WebSocket whose connector already resolved owner-only author
|
||||
# bindings. ``platform`` carries the UNDERLYING platform (not ``relay``), so
|
||||
# authz must key upstream trust off THIS flag. Excluded from to_dict/
|
||||
# from_dict so a peer can never forge or persist it.
|
||||
# Wire-INVISIBLE trust signal: delivered over the per-instance authenticated relay WebSocket
|
||||
# whose connector already resolved owner-only author bindings. ``platform`` carries the
|
||||
# UNDERLYING platform (not ``relay``), so authz must key upstream trust off THIS flag. Excluded
|
||||
# from to_dict/ from_dict so a peer can never forge or persist it.
|
||||
delivered_via_upstream_relay: bool = False
|
||||
|
||||
def __post_init__(self) -> None:
|
||||
@@ -350,9 +347,8 @@ def _discord_platform_notes(context: SessionContext) -> List[str]:
|
||||
else:
|
||||
lines.append(f" - Channel: `{src.chat_id}`")
|
||||
if src.message_id:
|
||||
# The volatile per-turn message id must stay OUT of this cached
|
||||
# block (it would bust the agent-cache signature every message);
|
||||
# run.py injects it into the user message instead.
|
||||
# The volatile per-turn message id must stay OUT of this cached block (it would bust the
|
||||
# agent-cache signature every message); run.py injects it into the user message instead.
|
||||
lines.append(
|
||||
" - Triggering message: provided per-turn in the incoming "
|
||||
"user message (use it as `message_id` for reply/react/pin)"
|
||||
@@ -467,9 +463,8 @@ def build_session_context_prompt(context: SessionContext, *, redact_pii: bool =
|
||||
"about other Matrix rooms or projects unless the user explicitly says so."
|
||||
)
|
||||
|
||||
# Shared multi-user sessions: never pin one user name in the system
|
||||
# prompt (changes per turn -> busts the prompt cache); sender names are
|
||||
# prefixed on each user message instead.
|
||||
# Shared multi-user sessions: never pin one user name in the system prompt (changes per turn ->
|
||||
# busts the prompt cache); sender names are prefixed on each user message instead.
|
||||
if context.shared_multi_user_session:
|
||||
session_label = "Multi-user thread" if src.thread_id else "Multi-user session"
|
||||
lines.append(
|
||||
@@ -565,8 +560,7 @@ class SessionEntry:
|
||||
# re-inject topic/channel skills. Distinct from was_auto_reset, which
|
||||
# fires the "expired due to inactivity" notice (wrong for a manual reset).
|
||||
is_fresh_reset: bool = False
|
||||
# Set by the expiry watcher after finalizing; persisted so restarts don't
|
||||
# re-run finalization.
|
||||
# Set by the expiry watcher after finalizing; persisted so restarts don't re-run finalization.
|
||||
expiry_finalized: bool = False
|
||||
# Next get_or_create_session() auto-resets (new session_id). Set by /stop
|
||||
# to break stuck-resume loops.
|
||||
@@ -761,10 +755,9 @@ def build_session_key(
|
||||
return ":".join(str(part) for part in dm_parts)
|
||||
|
||||
participant_id = _canonical_participant(source)
|
||||
# Discord auto-thread continuity: key a channel-initiating message on the
|
||||
# thread it WILL be delivered into (prospective_thread_id), and normalize
|
||||
# the chat_type slot to "thread" so in-thread follow-ups byte-match. A
|
||||
# real thread_id always wins.
|
||||
# Discord auto-thread continuity: key a channel-initiating message on the thread it WILL be
|
||||
# delivered into (prospective_thread_id), and normalize the chat_type slot to "thread" so
|
||||
# in-thread follow-ups byte-match. A real thread_id always wins.
|
||||
effective_thread_id = source.thread_id or source.prospective_thread_id
|
||||
chat_type_slot = source.chat_type
|
||||
if source.prospective_thread_id and not source.thread_id:
|
||||
@@ -882,30 +875,27 @@ class SessionStore(
|
||||
# Keep the legacy sessions.json mirror (disable via gateway.write_sessions_json).
|
||||
self._write_sessions_json = bool(getattr(config, "write_sessions_json", True))
|
||||
|
||||
# SQLite handles are cached per resolved path and resolved through the
|
||||
# ``_db`` property, never bound once here: a multiplexed gateway serves
|
||||
# every profile from ONE process and a handle frozen to the root home
|
||||
# would land every profile's rows in the root state.db. Priming the
|
||||
# current scope below keeps startup diagnostics at construction time.
|
||||
# SQLite handles are cached per resolved path and resolved through the ``_db`` property,
|
||||
# never bound once here: a multiplexed gateway serves every profile from ONE process and a
|
||||
# handle frozen to the root home would land every profile's rows in the root state.db.
|
||||
# Priming the current scope below keeps startup diagnostics at construction time.
|
||||
self._db_pinned = _DB_UNPINNED
|
||||
self._db_handles: Dict[Path, Any] = {}
|
||||
self._db_handles_lock = threading.Lock()
|
||||
# profile name -> HERMES_HOME; memoized so per-key store lookup is a
|
||||
# dict hit, not a profile-directory stat per append.
|
||||
self._profile_home_cache: Dict[str, Optional[Path]] = {}
|
||||
# session_id -> owning routing key for ids whose ownership is proven
|
||||
# but not yet published in ``_entries`` (compression child row is
|
||||
# written before its reroute is published).
|
||||
# session_id -> owning routing key for ids whose ownership is proven but not yet published
|
||||
# in ``_entries`` (compression child row is written before its reroute is published).
|
||||
self._session_owner_hints: Dict[str, str] = {}
|
||||
from gateway.session_db_recovery import RecoverableHandleCache
|
||||
|
||||
self._db_handle_cache = RecoverableHandleCache(
|
||||
handles=self._db_handles, lock=self._db_handles_lock,
|
||||
)
|
||||
# The routing index is one process-wide structure keyed by
|
||||
# ``agent:<profile>:…`` and needs exactly one home for its lifetime:
|
||||
# the gateway's own, captured before any profile scope exists
|
||||
# (see ``_routing_db``).
|
||||
# The routing index is one process-wide structure keyed by ``agent:<profile>:…`` and needs
|
||||
# exactly one home for its lifetime: the gateway's own, captured before any profile scope
|
||||
# exists (see ``_routing_db``).
|
||||
try:
|
||||
from hermes_constants import get_hermes_home
|
||||
|
||||
@@ -917,9 +907,8 @@ class SessionStore(
|
||||
def _lazy(self, name: str, factory):
|
||||
"""Return ``self.<name>``, creating it via *factory* when missing/None.
|
||||
|
||||
Suites build bare stores via ``object.__new__`` without running
|
||||
``__init__``; every optional lock/map is read through this so those
|
||||
instances still work.
|
||||
Suites build bare stores via ``object.__new__`` without running ``__init__``; every optional
|
||||
lock/map is read through this so those instances still work.
|
||||
"""
|
||||
value = getattr(self, name, None)
|
||||
if value is None:
|
||||
@@ -1077,9 +1066,8 @@ class SessionStore(
|
||||
stale_hit = checked and checks.is_stale
|
||||
reset_reason = checks.reset_reason if checked else None
|
||||
if stale_hit:
|
||||
# Stale routing self-heal: the entry points at a session
|
||||
# ALREADY ended in state.db. Drop it and fall through to
|
||||
# recovery (reopens agent_close / ws_orphan_reap rows,
|
||||
# Stale routing self-heal: the entry points at a session ALREADY ended in state.db.
|
||||
# Drop it and fall through to recovery (reopens agent_close / ws_orphan_reap rows,
|
||||
# fresh session for other end_reasons).
|
||||
logger.warning(
|
||||
"gateway.session: routing key %r -> %s is ended in "
|
||||
|
||||
@@ -1,9 +1,8 @@
|
||||
"""Session-scoped context variables for the Hermes gateway.
|
||||
|
||||
Replaces the old ``os.environ``-based ``HERMES_SESSION_*`` state with task-local
|
||||
``ContextVar``s (inherited by ``run_in_executor`` threads), so concurrently handled
|
||||
messages no longer clobber each other's routing ids. ``get_session_env`` is a drop-in
|
||||
for ``os.getenv``.
|
||||
Replaces the old ``os.environ``-based ``HERMES_SESSION_*`` state with task-local ``ContextVar``s
|
||||
(inherited by ``run_in_executor`` threads), so concurrently handled messages no longer clobber each
|
||||
other's routing ids. ``get_session_env`` is a drop-in for ``os.getenv``.
|
||||
"""
|
||||
|
||||
import os
|
||||
|
||||
@@ -43,10 +43,9 @@ def _parse_iso(value) -> Optional[datetime]:
|
||||
return None
|
||||
|
||||
|
||||
# Default auto-continue freshness window (1 hour): a restart-interrupted
|
||||
# session is only auto-resumed while within this window of when
|
||||
# ``resume_pending`` was marked. ``gateway/run.py`` bridges config.yaml
|
||||
# ``agent.gateway_auto_continue_freshness`` into the env var at startup.
|
||||
# Default auto-continue freshness window (1 hour): a restart-interrupted session is only
|
||||
# auto-resumed while within this window of when ``resume_pending`` was marked. ``gateway/run.py``
|
||||
# bridges config.yaml ``agent.gateway_auto_continue_freshness`` into the env var at startup.
|
||||
_AUTO_CONTINUE_FRESHNESS_SECS_DEFAULT = 60 * 60
|
||||
|
||||
|
||||
@@ -94,10 +93,9 @@ class SessionLifecycleMixin:
|
||||
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. Promotion only upgrades live/
|
||||
# agent_close rows; explicit boundaries are preserved.
|
||||
# Without a durable ``session_reset`` end_reason, later agent cleanup ends the row as
|
||||
# ``agent_close``, which stale-route recovery treats as resumable. Promotion only
|
||||
# upgrades live/ agent_close rows; explicit boundaries are preserved.
|
||||
_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)
|
||||
@@ -283,8 +281,7 @@ class SessionLifecycleMixin:
|
||||
started_at is None or (max_age_seconds > 0 and now - started_at > max_age)
|
||||
)
|
||||
except TypeError:
|
||||
# Mixed aware/naive timestamps: clear rather than risk an
|
||||
# unsafe old resume.
|
||||
# Mixed aware/naive timestamps: clear rather than risk an unsafe old resume.
|
||||
marker_is_stale = True
|
||||
|
||||
if not marker_is_stale and not entry.suspended:
|
||||
|
||||
@@ -1,6 +1,5 @@
|
||||
"""SessionStore storage plumbing: per-profile SessionDB handle resolution and the
|
||||
routing-index load/save paths (state.db gateway_routing primary, sessions.json
|
||||
legacy mirror).
|
||||
"""SessionStore storage plumbing: per-profile SessionDB handle resolution and the routing-index
|
||||
load/save paths (state.db gateway_routing primary, sessions.json legacy mirror).
|
||||
|
||||
Mixin split out of ``gateway/session.py``; bound onto ``SessionStore`` via the MRO.
|
||||
"""
|
||||
@@ -27,8 +26,7 @@ logger = logging.getLogger("gateway.session")
|
||||
# active scope" from a deliberate ``store._db = None`` (JSONL fallback).
|
||||
_DB_UNPINNED = object()
|
||||
|
||||
# Self-documenting sentinel written first into sessions.json; "_" keys are
|
||||
# skipped on load.
|
||||
# Self-documenting sentinel written first into sessions.json; "_" keys are skipped on load.
|
||||
_SESSIONS_JSON_README = (
|
||||
"LEGACY MIRROR of the gateway routing index (the primary copy "
|
||||
"lives in the gateway_routing table in ~/.hermes/state.db). "
|
||||
@@ -121,10 +119,9 @@ class SessionPersistenceMixin:
|
||||
def _named_profile_for_key(self, session_key: Optional[str]) -> Optional[str]:
|
||||
"""The non-default profile that owns *session_key*, or None.
|
||||
|
||||
None means the ambient store is authoritative (multiplexing off, or
|
||||
legacy ``agent:main`` namespace). It deliberately does NOT cover "that
|
||||
profile has no directory" — ownership and resolvability are separate
|
||||
questions that ``_db_for_key`` answers separately.
|
||||
None means the ambient store is authoritative (multiplexing off, or legacy ``agent:main``
|
||||
namespace). It deliberately does NOT cover "that profile has no directory" — ownership and
|
||||
resolvability are separate questions that ``_db_for_key`` answers separately.
|
||||
"""
|
||||
if not getattr(self.config, "multiplex_profiles", False):
|
||||
return None
|
||||
@@ -174,10 +171,9 @@ class SessionPersistenceMixin:
|
||||
return self._db
|
||||
home = self._profile_home_for_key(session_key)
|
||||
if home is None:
|
||||
# Named owner we cannot resolve (not provisioned yet, or lookup
|
||||
# failed). Falling back to the ambient store would split ONE
|
||||
# session identity across two physical stores — fail closed;
|
||||
# callers already handle a missing DB.
|
||||
# Named owner we cannot resolve (not provisioned yet, or lookup failed). Falling back to
|
||||
# the ambient store would split ONE session identity across two physical stores — fail
|
||||
# closed; callers already handle a missing DB.
|
||||
logger.warning(
|
||||
"gateway.session: profile %r has no resolvable home (key %r); "
|
||||
"refusing to fall back to the ambient store",
|
||||
@@ -400,9 +396,8 @@ class SessionPersistenceMixin:
|
||||
)
|
||||
return None
|
||||
|
||||
# Compression-ended parent with a newer live child for the same peer:
|
||||
# repoint instead of dropping, or queued/resume-pending work vanishes
|
||||
# until the next message.
|
||||
# Compression-ended parent with a newer live child for the same peer: repoint instead of
|
||||
# dropping, or queued/resume-pending work vanishes until the next message.
|
||||
if recovered_entry is not None and recovered_entry.session_id != entry.session_id:
|
||||
logger.warning(
|
||||
"gateway.session: repointing stale sessions.json entry "
|
||||
|
||||
@@ -251,9 +251,8 @@ class SessionRecoveryMixin:
|
||||
) -> tuple[Optional[SessionEntry], bool]:
|
||||
"""Find and gate a recoverable row -> (entry or None, migrated_legacy).
|
||||
|
||||
The legacy (pre-workspace) Slack key fallback lives here: exact-key
|
||||
lookup, claimed once per process; ``migrated_legacy`` tells the caller
|
||||
to rewrite the peer row to the scoped key.
|
||||
The legacy (pre-workspace) Slack key fallback lives here: exact-key lookup, claimed once per
|
||||
process; ``migrated_legacy`` tells the caller to rewrite the peer row to the scoped key.
|
||||
"""
|
||||
legacy_key = self._legacy_slack_session_key(source)
|
||||
recovered = self._find_gateway_session_row(
|
||||
|
||||
@@ -282,14 +282,12 @@ class SessionTranscriptMixin:
|
||||
return
|
||||
|
||||
if isinstance(exc, CompressionSessionClosedError):
|
||||
# Adopt only a different, still-live compression tip, else
|
||||
# fail closed.
|
||||
# Adopt only a different, still-live compression tip, else fail closed.
|
||||
_owner_key = self._owner_key_for_session_id(session_id)
|
||||
child_id = self._live_compression_child(session_id)
|
||||
if child_id:
|
||||
# Record the child's owner BEFORE writing to it (the
|
||||
# reroute is published only after the write succeeds
|
||||
# — load-bearing for backlog order).
|
||||
# Record the child's owner BEFORE writing to it (the reroute is published
|
||||
# only after the write succeeds — load-bearing for backlog order).
|
||||
if _owner_key:
|
||||
self._lazy("_session_owner_hints", dict)[child_id] = _owner_key
|
||||
try:
|
||||
|
||||
Reference in New Issue
Block a user