diff --git a/plugins/memory/honcho/__init__.py b/plugins/memory/honcho/__init__.py index 42564423c3..a774165973 100644 --- a/plugins/memory/honcho/__init__.py +++ b/plugins/memory/honcho/__init__.py @@ -634,11 +634,12 @@ class HonchoMemoryProvider(DialecticMixin, MemoryProvider): def sync_turn( self, user_content: str, assistant_content: str, *, session_id: str = "", - turn_author: Optional[Dict[str, Any]] = None, scope: Optional[str] = None, + turn_author: Optional[Dict[str, Any]] = None, ) -> None: """Record the conversation turn in Honcho (non-blocking), chunking messages that exceed the Honcho API limit. Honors saveMessages: false. ``turn_author`` names who wrote - the user side; the ``on_turn_start`` stash is the fallback for callers that never pass it.""" + the user side; the ``on_turn_start`` stash is the fallback for callers that never pass it. + A bot author's turn is written into that bot's own a2a session, never the human's.""" if not self._writes_enabled(): return if _is_internal_gateway_turn(user_content): @@ -656,11 +657,29 @@ class HonchoMemoryProvider(DialecticMixin, MemoryProvider): return author = turn_author if isinstance(turn_author, dict) else self._turn_author - # Resolved before the thread starts so a following turn cannot retag a queued write. - author_peer_id = self._manager.resolve_author_peer_id(self._session_key, author.get("id"), author.get("name")) + session_kwargs: dict[str, str] = {} + if author.get("is_bot"): + # A bot's turn never lands in the human's session: its own a2a session or nothing. + if not getattr(self._config, "a2a_sessions", True): + logger.debug("Honcho sync skipped a bot-authored turn because a2aSessions is off") + return + author_id = str(author.get("id") or "").strip() + if not author_id: + logger.debug("Honcho sync skipped a bot-authored turn that named no author id") + return + bot_peer_id = self._manager.resolve_author_peer_id(self._session_key, author_id, author.get("name")) + session_key = self._a2a_session_key(author_id) + # The bot is the a2a session's own user peer, so its messages need no per-message author. + if bot_peer_id: + session_kwargs["user_peer_id"] = bot_peer_id + author_peer_id = None + else: + session_key = self._session_key + # Resolved before the thread starts so a following turn cannot retag a queued write. + author_peer_id = self._manager.resolve_author_peer_id(session_key, author.get("id"), author.get("name")) def _sync(): - session = self._manager.get_or_create(self._session_key) + session = self._manager.get_or_create(session_key, **session_kwargs) for chunk in self._chunk_message(clean_user_content, msg_limit) if clean_user_content else (): session.add_message("user", chunk, author_peer_id=author_peer_id) for chunk in self._chunk_message(clean_assistant_content, msg_limit) if clean_assistant_content else (): @@ -672,6 +691,12 @@ class HonchoMemoryProvider(DialecticMixin, MemoryProvider): self._sync_thread.join(timeout=5.0) self._sync_thread = self._spawn_write(_sync, "honcho-sync", "Honcho sync_turn failed: %s") + def _a2a_session_key(self, author_id: str) -> str: + """Honcho session for one sender bot's DMs into this agent, stable across turns. Rooms already + get their own session by title, so only DMs reroute.""" + key = f"{self._session_key}:a2a:{re.sub(r'[^a-zA-Z0-9_-]', '-', author_id)}" + return HonchoClientConfig._enforce_session_id_limit(key, key) + @staticmethod def _spawn_write(fn: Callable[[], None], name: str, fail_msg: str) -> threading.Thread: """Run a Honcho write off-thread; failures are debug-logged, never raised into the turn.""" diff --git a/plugins/memory/honcho/client.py b/plugins/memory/honcho/client.py index 5d209b44ed..e600bb651a 100644 --- a/plugins/memory/honcho/client.py +++ b/plugins/memory/honcho/client.py @@ -321,6 +321,7 @@ def _behavior_fields(look: _HostLookup, explicitly_configured: bool) -> dict[str **_resolve_observation(observation_mode, look.pick("observation")), "session_strategy": look.pick("sessionStrategy", "per-directory"), "session_peer_prefix": look.pick_set("sessionPeerPrefix", False), + "a2a_sessions": look.flag("a2aSessions", default=True), } @@ -383,6 +384,8 @@ class HonchoClientConfig: # Session resolution session_strategy: str = "per-directory" session_peer_prefix: bool = False + # Bot-authored DMs write into their own session per sender bot. + a2a_sessions: bool = True sessions: dict[str, str] = field(default_factory=dict) raw: dict[str, Any] = field(default_factory=dict) # A hosts. block or explicit enabled flag, vs auto-enabled from a stray env key. diff --git a/plugins/memory/honcho/config_schema.py b/plugins/memory/honcho/config_schema.py index ebe35b6bd3..a30c7588f3 100644 --- a/plugins/memory/honcho/config_schema.py +++ b/plugins/memory/honcho/config_schema.py @@ -66,6 +66,9 @@ CONFIG_SCHEMA = ProviderConfigSchema( # — Session — _field("sessionPeerPrefix", "Session peer prefix", KIND_BOOL, "Prefix session peer names with the host.", default="false", group="Session"), + _field("a2aSessions", "Bot DM sessions", KIND_BOOL, + "Write DMs from other bots into their own Honcho session per sender. Off skips bot-authored turns.", + default="true", group="Session"), _field("sessions", "Session overrides", KIND_JSON, "Explicit session ID overrides keyed by resolver.", placeholder='{"key": "session-id"}', group="Session", scope="root"), # — Message writing — diff --git a/plugins/memory/honcho/session.py b/plugins/memory/honcho/session.py index d2f05cfb6d..b0185c2b39 100644 --- a/plugins/memory/honcho/session.py +++ b/plugins/memory/honcho/session.py @@ -213,8 +213,10 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi honcho_session = self._authed_call("session setup", lambda: self._sdk_session(session_id)) return honcho_session, existing_messages - def get_or_create(self, key: str) -> HonchoSession: - """Get an existing session or create a new one for ``key`` (usually channel:chat_id).""" + def get_or_create(self, key: str, *, user_peer_id: str | None = None) -> HonchoSession: + """Get an existing session or create a new one for ``key`` (usually channel:chat_id). + ``user_peer_id`` replaces the resolved user peer when the session's participant is not the + runtime user, e.g. the sender bot of an a2a session.""" with self._cache_lock: if key in self._cache: logger.debug("Local session cache hit: %s", key) @@ -224,7 +226,7 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi # bots scope memory per user; config can alias/prefix it, or pinPeerName pins all # identities to peerName for single-user deployments (see _resolve_user_peer_id). # Determine peer IDs — no lock needed (read-only, no shared state mutation). See #14984. - user_peer_id = self._resolve_user_peer_id(key) + user_peer_id = user_peer_id or self._resolve_user_peer_id(key) assistant_peer_id = self._sanitize_id(self._config.ai_peer if self._config else "hermes-assistant") # All expensive I/O outside the lock — Honcho's persistence is source of truth. diff --git a/tests/honcho_plugin/test_a2a_sessions.py b/tests/honcho_plugin/test_a2a_sessions.py new file mode 100644 index 0000000000..619f07a50e --- /dev/null +++ b/tests/honcho_plugin/test_a2a_sessions.py @@ -0,0 +1,181 @@ +"""Bot-authored DMs write into their own Honcho session. + +The turn author carries ``is_bot`` for a relayed DM. With ``a2aSessions`` on +(the default) ``sync_turn`` writes the whole turn into ``:a2a:`` +with the sender bot as that session's user peer. The human's session never +receives a bot's turn: with the flag off the turn is skipped. +""" + +import json +from types import SimpleNamespace +from unittest.mock import MagicMock + +from plugins.memory.honcho import HonchoMemoryProvider +from plugins.memory.honcho.client import HonchoClientConfig +from plugins.memory.honcho.session import HonchoSessionManager + +BOT_AUTHOR = {"id": "bot:coder", "name": "coder", "is_bot": True} +HUMAN_AUTHOR = {"id": "111222", "name": "Alice", "is_bot": False} + + +def _provider(a2a_sessions: bool = True) -> HonchoMemoryProvider: + provider = HonchoMemoryProvider() + provider._session_key = "Bot-Chat" + provider._manager = MagicMock() + provider._manager.get_or_create.return_value = MagicMock() + provider._cron_skipped = False + provider._session_initialized = True + provider._config = SimpleNamespace(message_max_chars=25000, a2a_sessions=a2a_sessions) + return provider + + +def _sync(provider: HonchoMemoryProvider, **kwargs) -> None: + provider.sync_turn("Message from coder: hi", "hello coder", **kwargs) + if provider._sync_thread: + provider._sync_thread.join(timeout=5) + + +class TestA2aRouting: + def test_bot_turn_lands_in_its_own_session(self): + provider = _provider() + provider._manager.resolve_author_peer_id.return_value = "coder" + + _sync(provider, turn_author=BOT_AUTHOR) + + provider._manager.get_or_create.assert_called_once_with("Bot-Chat:a2a:bot-coder", user_peer_id="coder") + session = provider._manager.get_or_create.return_value + roles = [c[0][0] for c in session.add_message.call_args_list] + assert roles == ["user", "assistant"] + # The bot is the session's user peer; no per-message author is attached. + assert session.add_message.call_args_list[0][1]["author_peer_id"] is None + + def test_bot_turn_never_touches_the_human_session(self): + provider = _provider() + provider._manager.resolve_author_peer_id.return_value = "coder" + + _sync(provider, turn_author=BOT_AUTHOR) + + keys = [c[0][0] for c in provider._manager.get_or_create.call_args_list] + assert "Bot-Chat" not in keys + + def test_session_key_is_stable_across_turns(self): + provider = _provider() + provider._manager.resolve_author_peer_id.return_value = "coder" + + _sync(provider, turn_author=BOT_AUTHOR) + _sync(provider, turn_author=BOT_AUTHOR) + + keys = {c[0][0] for c in provider._manager.get_or_create.call_args_list} + assert keys == {"Bot-Chat:a2a:bot-coder"} + + def test_two_bots_get_two_sessions(self): + provider = _provider() + provider._manager.resolve_author_peer_id.side_effect = lambda key, author_id, name=None: author_id[4:] + + _sync(provider, turn_author=BOT_AUTHOR) + _sync(provider, turn_author={"id": "bot:writer", "name": "writer", "is_bot": True}) + + keys = [c[0][0] for c in provider._manager.get_or_create.call_args_list] + assert keys == ["Bot-Chat:a2a:bot-coder", "Bot-Chat:a2a:bot-writer"] + + def test_bot_author_without_an_id_is_skipped(self): + provider = _provider() + + _sync(provider, turn_author={"id": None, "name": "mystery", "is_bot": True}) + + provider._manager.get_or_create.assert_not_called() + + def test_collapsed_bot_peer_keeps_the_default_user_peer(self): + """pinUserPeer returns no author peer; the a2a session then falls back to the resolved peer.""" + provider = _provider() + provider._manager.resolve_author_peer_id.return_value = None + + _sync(provider, turn_author=BOT_AUTHOR) + + provider._manager.get_or_create.assert_called_once_with("Bot-Chat:a2a:bot-coder") + + def test_long_session_key_stays_within_the_honcho_limit(self): + provider = _provider() + provider._session_key = "x" * 95 + provider._manager.resolve_author_peer_id.return_value = "coder" + + _sync(provider, turn_author=BOT_AUTHOR) + + key = provider._manager.get_or_create.call_args[0][0] + assert len(key) <= 100 + assert key == provider._a2a_session_key("bot:coder") + + +class TestFlagOff: + def test_bot_turn_is_skipped(self): + provider = _provider(a2a_sessions=False) + + _sync(provider, turn_author=BOT_AUTHOR) + + provider._manager.get_or_create.assert_not_called() + provider._manager.resolve_author_peer_id.assert_not_called() + + def test_human_turn_still_writes(self): + provider = _provider(a2a_sessions=False) + provider._manager.resolve_author_peer_id.return_value = "alice" + + _sync(provider, turn_author=HUMAN_AUTHOR) + + provider._manager.get_or_create.assert_called_once_with("Bot-Chat") + + +class TestHumanTurnUnchanged: + def test_human_turn_writes_into_the_session_under_the_author(self): + provider = _provider() + provider._manager.resolve_author_peer_id.return_value = "alice" + + _sync(provider, turn_author=HUMAN_AUTHOR) + + provider._manager.get_or_create.assert_called_once_with("Bot-Chat") + session = provider._manager.get_or_create.return_value + assert session.add_message.call_args_list[0][1]["author_peer_id"] == "alice" + + def test_human_turn_without_author_keeps_the_session_peer(self): + provider = _provider() + provider._manager.resolve_author_peer_id.return_value = None + + _sync(provider) + + provider._manager.get_or_create.assert_called_once_with("Bot-Chat") + + +class TestManagerUserPeerOverride: + def test_get_or_create_uses_the_override_as_user_peer(self): + mgr = HonchoSessionManager(honcho=MagicMock(), config=HonchoClientConfig(api_key="k", peer_name="eri", ai_peer="hermes"), + runtime_user_peer_name="7654321") + mgr._get_or_create_peer = MagicMock(side_effect=lambda pid: MagicMock(name=f"peer:{pid}")) + mgr._get_or_create_honcho_session = MagicMock(return_value=(MagicMock(), [])) + + session = mgr.get_or_create("Bot-Chat:a2a:bot-coder", user_peer_id="coder") + + assert session.user_peer_id == "coder" + assert session.assistant_peer_id == "hermes" + joined = [c[0][0] for c in mgr._get_or_create_peer.call_args_list] + assert "7654321" not in joined + + +class TestConfigFlag: + def _config(self, tmp_path, monkeypatch, raw: dict) -> HonchoClientConfig: + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + path = tmp_path / "honcho.json" + path.write_text(json.dumps({"apiKey": "k", **raw})) + return HonchoClientConfig.from_global_config(config_path=path) + + def test_defaults_on(self, tmp_path, monkeypatch): + assert self._config(tmp_path, monkeypatch, {}).a2a_sessions is True + + def test_root_flag(self, tmp_path, monkeypatch): + assert self._config(tmp_path, monkeypatch, {"a2aSessions": False}).a2a_sessions is False + + def test_host_block_wins_over_root(self, tmp_path, monkeypatch): + cfg = self._config(tmp_path, monkeypatch, {"a2aSessions": True, "hosts": {"hermes": {"a2aSessions": False}}}) + assert cfg.a2a_sessions is False + + def test_root_applies_when_host_block_is_silent(self, tmp_path, monkeypatch): + cfg = self._config(tmp_path, monkeypatch, {"a2aSessions": False, "hosts": {"hermes": {"peerName": "eri"}}}) + assert cfg.a2a_sessions is False diff --git a/tests/plugins/memory/test_honcho_config_schema.py b/tests/plugins/memory/test_honcho_config_schema.py index 8ed43e8ac3..a8434f088b 100644 --- a/tests/plugins/memory/test_honcho_config_schema.py +++ b/tests/plugins/memory/test_honcho_config_schema.py @@ -50,6 +50,8 @@ def test_declares_the_new_field_kinds(): by_key = {f.key: f for f in provider.fields} assert by_key["saveMessages"].kind == KIND_BOOL + assert by_key["a2aSessions"].kind == KIND_BOOL + assert by_key["a2aSessions"].default == "true" assert by_key["dialecticMaxChars"].kind == KIND_NUMBER assert by_key["userPeerAliases"].kind == KIND_JSON assert by_key["recallMode"].allowed_values() == {"hybrid", "context", "tools"}