Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions gateway/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,20 @@ gateway under the backend, and do NOT "fix" update locks by widening the tree-ki
| shared bot → profile that owns its own bot | routed | receiving adapter | receiving adapter (conversation continuity; #70625's "routed bot" reading was not adopted) |
| secondary-owned bot → `default` (`bot_profile`) | default (`agent:main`) | receiving adapter | receiving adapter |
| restored / hand-built, no live provenance | stored `source.profile` | **`None`** | unique owner of `(platform, runtime)`; a disconnected secondary → `None`, never the default bot |
- **Identity survives the process.** `SessionEntry.transport_profile` (routing index +
`sessions.transport_profile`, nullable, reconciled by `SCHEMA_SQL`) persists the receiving bot
next to the key namespace; the namespace says where a lane RUNS, the column says which bot may
DELIVER to it. Anything reviving a session from durable state (auto-resume, heartbeat restore,
plugin injection, background-process events) reads `entry.origin` through
`authz_mixin.py::_restored_source`, which re-pins a `RoutingIdentity(transport=None)` via
`session_identity.restore_identity`; `_delivery_adapter_for` then delivers through that bot's
adapter or nothing (never the default bot by heuristic). Rows without the column (pre-PR-5)
keep the `_is_shared_bot_satellite` fallback. Deferred callbacks capture the identity/home at
command time (`/model` picker); `_run_in_executor_with_context` carries the scope over thread
hops. Relay: `_with_scope` echoes the routed `profile` on every outbound frame and `follow_up`
derives it from the key namespace so the connector stamps it on the next `passthrough_forward`.
Ambient `get_active_profile_name()` reads in `gateway/` are boot-only and marked
`# launch profile, pre-identity`; a path with an identity reads `identity.runtime_profile`.
- **Token locks.** An adapter that connects with a unique credential (bot token, API key) calls
`acquire_scoped_lock()` from `gateway.status` in `connect()`/`start()` and `release_scoped_lock()`
in `disconnect()`/`stop()`, so two profiles cannot share one credential. Canonical:
Expand Down
32 changes: 31 additions & 1 deletion gateway/authz_mixin.py
Original file line number Diff line number Diff line change
Expand Up @@ -300,6 +300,17 @@ def _delivery_adapter_for(self, source: Optional[SessionSource]):
adapter = self._intake_adapter_for(source)
if adapter is not None:
return adapter
# A pinned identity NAMES the receiving bot (live or restored from ``transport_profile``).
# If that bot has no adapter right now it is offline: fail closed rather than fall through to
# the runtime profile's bot — that fallthrough is the "restored lane answers from the wrong
# bot" row the identity exists to close. Only an identity-less source (legacy row, bare
# fixture) uses the unique-owner heuristic below.
from gateway.session_identity import identity_of
identity = identity_of(source)
if identity is not None and identity.multiplexed and not identity.transport_inferred:
return None
# No identity, or one whose transport was only inferred (hand-built source, pre-column row):
# the unique owner of ``(platform, runtime_profile)`` delivers.
# ``getattr``: test fixtures build bare SimpleNamespace sources without ``profile``.
return self._authorization_adapter(getattr(source, "platform", None), getattr(source, "profile", None))

Expand Down Expand Up @@ -360,7 +371,26 @@ def _authorization_home_for_source(self, source: SessionSource):
def _adapter_profile_for_source(self, source: SessionSource) -> Optional[str]:
"""Resolve the transport-owning profile for adapter policy lookups."""
owner = self._transport_owner(source)
return owner[1] if owner is not None else getattr(source, "profile", None)
if owner is not None:
return owner[1]
from gateway.session_identity import identity_of
identity = identity_of(source)
if identity is not None and identity.multiplexed:
return None if identity.transport_profile == "default" else identity.transport_profile
return getattr(source, "profile", None)

def _restored_source(self, entry) -> Optional[SessionSource]:
"""``entry.origin`` with its identity re-pinned from the routing entry's persisted
``transport_profile`` (no live adapter: the restored row of the transport matrix). Every path
that revives a session from durable state — auto-resume, heartbeat restore, plugin injection,
background-process events — reads the origin through here, so the receiving bot decides
delivery and authorization after a restart, not the runtime profile's heuristics."""
source = getattr(entry, "origin", None)
if source is None:
return None
from gateway.session_identity import restore_identity
restore_identity(source, runner=self, transport_profile=getattr(entry, "transport_profile", None))
return source

def _adapter_flag(self, platform, name: str, profile) -> bool:
"""Adapter-declared boolean, False when unknown. ``authorization_is_upstream`` (relay: a trusted
Expand Down
2 changes: 1 addition & 1 deletion gateway/platforms/api_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -1308,7 +1308,7 @@ def _resolve_model_name(explicit: str) -> str:
profile_name = ""
with suppress(Exception):
from hermes_cli.profiles import get_active_profile_name
profile = get_active_profile_name()
profile = get_active_profile_name() # launch profile, pre-identity (advertised model name)
if profile and profile not in {"default", "custom"}:
profile_name = profile
return resolve_effective_model(explicit, profile_name, "hermes-agent")
Expand Down
31 changes: 29 additions & 2 deletions gateway/relay/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,17 @@ def _event_ids(event) -> Tuple[Optional[str], Optional[str]]:
return message_id, getattr(event.source, "chat_id", None)


def _profile_from_session_key(session_key: str) -> Optional[str]:
"""Named profile encoded in an ``agent:<ns>:...`` session key; None for the legacy ``agent:main``
namespace (single-profile gateway) so the wire frame stays byte-identical there."""
parts = (session_key or "").split(":")
if len(parts) < 2 or parts[0] != "agent" or not parts[1]:
return None
from gateway.session import profile_from_session_key_namespace
profile = profile_from_session_key_namespace(parts[1])
return None if profile == "default" else profile


class RelayAdapter(BasePlatformAdapter):
"""Generic relay adapter advertising a connector-negotiated capability profile."""

Expand Down Expand Up @@ -125,6 +136,10 @@ def __init__(
# platforms on one WS and a reply must egress through the platform the
# inbound came from. Empty for a single-platform gateway (connector default).
self._platform_by_chat: Dict[str, str] = {}
# chat_id -> Hermes profile the connector routed the inbound to (multiplex mode). Echoed
# on every outbound frame's metadata so the connector can stamp the SAME profile on the
# next passthrough_forward for that chat; empty on a single-profile gateway.
self._profile_by_chat: Dict[str, str] = {}
# Chats the connector has refused (see the terminal-decline latch).
# chat_id -> (thread_id, initial_name) of the auto-thread the CONNECTOR
# created for our latest send; read by the semantic thread-rename lane.
Expand Down Expand Up @@ -1042,6 +1057,7 @@ def _capture_scope(self, event) -> None:
for attr, cache in (
("user_id", self._dm_user_by_chat), ("scope_id", self._scope_by_chat),
("chat_type", self._chat_type_by_chat),
("profile", self.__dict__.setdefault("_profile_by_chat", {})),
):
value = getattr(src, attr, None)
if value:
Expand All @@ -1059,7 +1075,11 @@ def _with_scope(self, chat_id: str, metadata: Optional[Dict[str, Any]]) -> Dict[
first and only falls back to user_id on a route miss, so carrying both never
overrides routing-table resolution."""
meta: Dict[str, Any] = dict(metadata or {})
for key, cache in (("scope_id", self._scope_by_chat), ("user_id", self._dm_user_by_chat)):
# ``getattr``: relay tests build bare adapters via ``__new__`` without ``__init__``.
for key, cache in (
("scope_id", self._scope_by_chat), ("user_id", self._dm_user_by_chat),
("profile", getattr(self, "_profile_by_chat", {})),
):
if not meta.get(key):
value = cache.get(str(chat_id))
if value:
Expand Down Expand Up @@ -1698,13 +1718,20 @@ async def send_follow_up(
# default routes it.
prefix = kind.split(".", 1)[0] if kind and "." in kind else None
follow_up_platform = prefix if prefix and self.fronts_platform(prefix) else None
follow_up_metadata = dict(metadata or {})
# The session key names the profile namespace the interaction ran under; carry it so the
# connector's next passthrough_forward for this interaction routes to the same profile.
if not follow_up_metadata.get("profile"):
profile = _profile_from_session_key(session_key)
if profile:
follow_up_metadata["profile"] = profile
result = await self._transport.send_follow_up(
{
"op": "follow_up",
"session_key": session_key,
"kind": kind,
"content": content,
"metadata": metadata or {},
"metadata": follow_up_metadata,
},
platform=follow_up_platform,
)
Expand Down
11 changes: 8 additions & 3 deletions gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -1613,7 +1613,7 @@ def _cron_tick_profile_homes(config: object) -> list[tuple[str, "Path"]]:
from hermes_cli.profiles import get_active_profile_name, get_profile_dir

homes = _multiplex_profile_homes(config)
active = get_active_profile_name() or "default"
active = get_active_profile_name() or "default" # launch profile, pre-identity (ticker boot)
if any(name == active for name, _home in homes):
return homes
try:
Expand Down Expand Up @@ -3841,9 +3841,14 @@ def _session_key_for_source(self, source: SessionSource) -> str:
pass
config = getattr(self, "config", None)
# Mirror SessionStore._resolve_profile_for_key so this fallback yields the primary path's
# namespace: None (legacy agent:main) unless multiplexing is on, then the active profile.
# namespace: None (legacy agent:main) unless multiplexing is on, then the pinned identity's
# runtime profile, the source stamp, or the active profile.
from gateway.session_identity import identity_of
identity = identity_of(source)
_profile = None
if getattr(config, "multiplex_profiles", False):
if identity is not None:
_profile = identity.session_key_profile
elif getattr(config, "multiplex_profiles", False):
if source.profile:
_profile = source.profile
else:
Expand Down
2 changes: 1 addition & 1 deletion gateway/run_adapters.py
Original file line number Diff line number Diff line change
Expand Up @@ -853,7 +853,7 @@ async def _start_secondary_profile_adapters(self) -> int:
from hermes_cli.profiles import get_active_profile_name
except Exception:
return 0
active = get_active_profile_name() or "default"
active = get_active_profile_name() or "default" # launch profile, pre-identity (adapter boot)
connected = 0
claimed = self._primary_resource_claims(active)
profile_homes = _multiplex_profile_homes(self.config)
Expand Down
5 changes: 3 additions & 2 deletions gateway/run_heartbeat_restore.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,10 +48,11 @@ def scan():
if entry.origin is None or not entry.session_id or entry.suspended:
continue
try:
with runner._profile_scope_for_source(entry.origin):
source = runner._restored_source(entry)
with runner._profile_scope_for_source(source):
manager = HeartbeatManager(entry.session_id)
if manager.is_active():
restored.append((entry.session_key, entry.origin, entry.session_id))
restored.append((entry.session_key, source, entry.session_id))
except Exception:
logger.debug("heartbeat restore for %s failed", entry.session_key, exc_info=True)
return restored
Expand Down
3 changes: 2 additions & 1 deletion gateway/run_inbound.py
Original file line number Diff line number Diff line change
Expand Up @@ -1857,7 +1857,8 @@ def _accepting() -> bool:
if entry is None or entry.origin is None or not _accepting():
return False

source = dataclasses.replace(entry.origin)
from gateway.session_identity import replace_source
source = replace_source(self._restored_source(entry))
try:
authorized = self._is_user_authorized_for_source(source, allow_adapter_delegation=False)
except Exception:
Expand Down
2 changes: 1 addition & 1 deletion gateway/run_notifications.py
Original file line number Diff line number Diff line change
Expand Up @@ -969,7 +969,7 @@ def _build_process_event_source(self, evt: dict):
self.session_store._ensure_loaded()
entry = self.session_store._entries.get(session_key)
if entry and getattr(entry, "origin", None):
return entry.origin
return self._restored_source(entry)
except Exception as exc:
logger.debug("Synthetic process-event session-store lookup failed for %s: %s", session_key, exc)
cached_source = self._get_cached_session_source(session_key)
Expand Down
4 changes: 2 additions & 2 deletions gateway/run_startup.py
Original file line number Diff line number Diff line change
Expand Up @@ -566,7 +566,7 @@ def _schedule_resume_pending_sessions(self, platform=None) -> int:
# Already being resumed (e.g. scheduled at startup, still in-flight) — no second turn.
if self._is_session_running(entry.session_key):
continue
source = entry.origin
source = self._restored_source(entry)
adapter = self._delivery_adapter_for(source)
if adapter is None:
logger.debug(
Expand Down Expand Up @@ -816,7 +816,7 @@ def _start_log_startup_environment(self) -> None:
)
with suppress(Exception):
from hermes_cli.profiles import get_active_profile_name
_profile = get_active_profile_name()
_profile = get_active_profile_name() # launch profile, pre-identity (boot log)
if _profile and _profile != "default":
logger.info("Active profile: %s", _profile)
_write_runtime_status_quiet(gateway_state="starting", exit_reason=None, clear_profile_platforms=True)
Expand Down
21 changes: 17 additions & 4 deletions gateway/session.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@

from .config import Platform, GatewayConfig, HomeChannel
from .whatsapp_identity import canonical_whatsapp_identifier
from gateway.session_identity import transport_profile_of
from gateway.session_persistence import SessionPersistenceMixin, _DB_UNPINNED
from gateway.session_recovery import SessionRecoveryMixin
from gateway.session_lifecycle import SessionLifecycleMixin, _iso, _new_session_id, _now, _parse_iso
Expand Down Expand Up @@ -519,6 +520,10 @@ class SessionEntry:
# Session-scoped /model override (model/provider/base_url ONLY — never credentials, see
# sanitize_model_override). Persisted so a restart keeps the chosen model.
model_override: Optional[Dict[str, str]] = None
# Profile owning the bot that received this lane's traffic (``RoutingIdentity.transport_profile``,
# "default" spelled out). The key namespace only says where the turn RUNS; after a restart this is
# what says which bot may deliver to it. None = unknown (row predates the field, or standalone).
transport_profile: Optional[str] = None

# Fields (de)serialized verbatim, in wire order (``from_dict`` reads them with
# ``data.get(name, <dataclass default>)``), split around the three ISO-datetime/token keys.
Expand Down Expand Up @@ -548,6 +553,8 @@ def to_dict(self) -> Dict[str, Any]:
if self.model_override:
# Defence-in-depth against an unsanitized dict stored directly.
result["model_override"] = sanitize_model_override(self.model_override)
if self.transport_profile:
result["transport_profile"] = self.transport_profile
if self.origin:
result["origin"] = self.origin.to_dict()
return result
Expand Down Expand Up @@ -578,6 +585,7 @@ def from_dict(cls, data: Dict[str, Any]) -> "SessionEntry":
defaults = {f.name: f.default for f in fields(cls)}
plain = {n: data.get(n, defaults[n]) for n in cls._PLAIN_FIELDS + cls._RESET_FIELDS}
plain["expiry_finalized"] = data.get("expiry_finalized", data.get("memory_flushed", False))
transport_profile = data.get("transport_profile")
return cls(
session_key=session_key, session_id=session_id,
created_at=datetime.fromisoformat(data["created_at"]),
Expand All @@ -586,7 +594,9 @@ def from_dict(cls, data: Dict[str, Any]) -> "SessionEntry":
chat_type=data.get("chat_type", "dm"), metadata=dict(data.get("metadata") or {}),
last_resume_marked_at=_parse_iso(data.get("last_resume_marked_at")),
active_turn_token=token, active_turn_started_at=started_at,
model_override=sanitize_model_override(data.get("model_override")), **plain,
model_override=sanitize_model_override(data.get("model_override")),
transport_profile=transport_profile if isinstance(transport_profile, str) and transport_profile else None,
**plain,
)


Expand Down Expand Up @@ -1013,7 +1023,7 @@ def _route_create(
origin=source, display_name=source.chat_name, platform=source.platform,
chat_type=source.chat_type, was_auto_reset=decision.reset_reason is not None,
auto_reset_reason=decision.reset_reason, reset_had_activity=decision.reset_had_activity,
prev_session_id=decision.prev_session_id,
prev_session_id=decision.prev_session_id, transport_profile=transport_profile_of(source),
)
with self._lock:
current = self._entries.get(session_key)
Expand Down Expand Up @@ -1044,9 +1054,11 @@ def update_session(
entry.last_prompt_tokens = last_prompt_tokens
# Snapshot peer fields under _lock so a concurrent reset/heal cannot tear the row.
peer_sid, peer_origin, peer_name = entry.session_id, entry.origin, entry.display_name
peer_transport = entry.transport_profile
# Metadata-only: single-row UPSERT, outside ``_lock``.
self._save_entry(session_key)
self._record_gateway_session_peer(peer_sid, session_key, peer_origin, display_name=peer_name)
self._record_gateway_session_peer(
peer_sid, session_key, peer_origin, display_name=peer_name, transport_profile=peer_transport)

def get_session_metadata(self, session_key: str, key: str, default: Any = None) -> Any:
"""Return a metadata value stored on a live session entry."""
Expand Down Expand Up @@ -1117,7 +1129,7 @@ def _replace_route_locked(self, session_key, old_entry, session_id, now, **field
new_entry = SessionEntry(
session_key=session_key, session_id=session_id, created_at=now, updated_at=now,
origin=old_entry.origin, platform=old_entry.platform, chat_type=old_entry.chat_type,
**fields,
transport_profile=old_entry.transport_profile, **fields,
)
self._entries[session_key] = new_entry
self._save()
Expand Down Expand Up @@ -1212,6 +1224,7 @@ def switch_session(
self._record_gateway_session_peer(
target_session_id, session_key, new_entry.origin,
display_name=new_entry.display_name, include_compression_ancestors=True,
transport_profile=new_entry.transport_profile,
)
return new_entry

Expand Down
Loading
Loading