diff --git a/plugins/memory/honcho/__init__.py b/plugins/memory/honcho/__init__.py index dc6662fd05..b863c6c491 100644 --- a/plugins/memory/honcho/__init__.py +++ b/plugins/memory/honcho/__init__.py @@ -120,6 +120,8 @@ class HonchoMemoryProvider(DialecticMixin, MemoryProvider): # Recall cadence state (overwritten from config in initialize()). self._turn_count = 0 + # Author of the turn in flight, refreshed by on_turn_start. + self._turn_author: dict[str, Any] = {} self._query_rewrite_enabled = False self._injection_frequency = "every-turn" # or "first-turn" self._context_cadence = 1 # minimum turns between context API calls @@ -551,9 +553,12 @@ class HonchoMemoryProvider(DialecticMixin, MemoryProvider): _is_trivial_prompt = staticmethod(is_trivial_prompt) def on_turn_start(self, turn_number: int, message: str, **kwargs) -> None: - """Track turn count for cadence and injection_frequency logic.""" + """Track turn count for cadence, and record who wrote this turn: a shared session carries + several participants, and the peer resolved at session init only names whoever opened it.""" self._recall_generation = object() self._turn_count = turn_number + self._turn_author = {"id": kwargs.get("author_id") or None, "name": kwargs.get("author_name") or None, + "is_bot": bool(kwargs.get("author_is_bot"))} def on_session_switch(self, new_session_id: str, **kwargs) -> None: """Discard in-flight recall even when the configured backend session is pinned.""" @@ -615,11 +620,16 @@ class HonchoMemoryProvider(DialecticMixin, MemoryProvider): if not clean_user_content and not clean_assistant_content: return + # 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, self._turn_author.get("id"), self._turn_author.get("name")) + def _sync(): session = self._manager.get_or_create(self._session_key) - for role, content in (("user", clean_user_content), ("assistant", clean_assistant_content)): - for chunk in self._chunk_message(content, msg_limit) if content else (): - session.add_message(role, chunk) + 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 (): + session.add_message("assistant", chunk) # save() (not _flush_session) so writeFrequency batching is honored. self._manager.save(session) diff --git a/plugins/memory/honcho/session.py b/plugins/memory/honcho/session.py index 0d94308124..d2f05cfb6d 100644 --- a/plugins/memory/honcho/session.py +++ b/plugins/memory/honcho/session.py @@ -63,6 +63,8 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi self._cache: dict[str, HonchoSession] = {} self._cache_lock = threading.RLock() self._peers_cache: dict[str, Any] = {} + # honcho_session_id -> author peer IDs already joined to that session. + self._joined_author_peers: dict[str, set[str]] = {} self._sessions_cache: dict[str, Any] = {} # Bumped (under _cache_lock) whenever _force_reauth rebuilds the client, so an # in-flight resolver never stores an object bound to the discarded client. @@ -245,6 +247,26 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi # ----- Writes ----- + def _author_peer_for_session(self, honcho_session: Any, honcho_session_id: str, author_peer_id: str) -> Any: + """Return the author's peer, joining it to the session on first sight. A shared session's + roster is open (people and agents arrive later), so peers join when they first write; joins + are remembered per session so this costs one API call per author.""" + peer = self._get_or_create_peer(author_peer_id) + with self._cache_lock: + if author_peer_id in self._joined_author_peers.setdefault(honcho_session_id, set()): + return peer + try: + from honcho.session import SessionPeerConfig + config = SessionPeerConfig(observe_me=self._user_observe_me, observe_others=self._user_observe_others) + honcho_session.add_peers([(peer, config)]) + except Exception as e: + # The write still lands under the right peer; only the membership (observe config) is missing. + logger.debug("Honcho author peer join failed for %s: %s", author_peer_id, e) + return peer + with self._cache_lock: + self._joined_author_peers.setdefault(honcho_session_id, set()).add(author_peer_id) + return peer + def _flush_session(self, session: HonchoSession) -> bool: """Write unsynced messages to Honcho synchronously.""" new_messages = [m for m in session.messages if not m.get("_synced")] @@ -258,7 +280,15 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi honcho_session = self._sessions_cache.get(session.honcho_session_id) if honcho_session is None: honcho_session, _ = self._get_or_create_honcho_session(session.honcho_session_id, user_peer, assistant_peer) - honcho_messages = [(user_peer if m["role"] == "user" else assistant_peer).message(m["content"]) for m in new_messages] + honcho_messages = [] + for m in new_messages: + if m["role"] != "user": + honcho_messages.append(assistant_peer.message(m["content"])) + continue + author_peer_id = m.get("author_peer_id") + peer = (self._author_peer_for_session(honcho_session, session.honcho_session_id, author_peer_id) + if author_peer_id else user_peer) + honcho_messages.append(peer.message(m["content"])) honcho_session.add_messages(honcho_messages) return len(honcho_messages) diff --git a/plugins/memory/honcho/session_peers.py b/plugins/memory/honcho/session_peers.py index cf28865c7d..40ba510d50 100644 --- a/plugins/memory/honcho/session_peers.py +++ b/plugins/memory/honcho/session_peers.py @@ -100,3 +100,28 @@ class SessionPeersMixin: if self._ai_observe_others: return session.assistant_peer_id, target_peer_id return target_peer_id, None + + def _peer_id_for_runtime_id(self, runtime_id: str) -> str: + """Map one gateway runtime identity onto its Honcho peer ID with the same alias-then-prefix + order ``_resolve_user_peer_id`` applies, so an aliased account lands on its peer on any turn.""" + aliases = getattr(self._config, "user_peer_aliases", {}) if self._config else {} + if isinstance(aliases, dict): + alias = aliases.get(runtime_id) + if isinstance(alias, str) and alias.strip(): + return self._sanitize_id(alias.strip()) + prefix = getattr(self._config, "runtime_peer_prefix", "") if self._config else "" + prefix = prefix.strip() if isinstance(prefix, str) else "" + return self._generated_runtime_peer_id(prefix, runtime_id) if prefix else self._sanitize_id(runtime_id) + + def resolve_author_peer_id(self, key: str, author_id: str | None, author_name: str | None = None) -> str | None: + """Peer ID for the turn's author, or None to keep the session's peer: no author named, the + author IS the session's peer, or ``pinPeerName`` collapsing identities by operator request. + ``author_name`` is a display name (attacker-influenceable), so it never becomes a peer ID.""" + runtime_id = str(author_id).strip() if author_id else "" + if not runtime_id: + return None + if self._config is not None and bool(getattr(self._config, "peer_name", None)) \ + and getattr(self._config, "pin_peer_name", False) is True: + return None + peer_id = self._peer_id_for_runtime_id(runtime_id) + return None if peer_id == self._resolve_user_peer_id(key) else peer_id diff --git a/tests/honcho_plugin/test_turn_author_peers.py b/tests/honcho_plugin/test_turn_author_peers.py new file mode 100644 index 0000000000..57e0bad499 --- /dev/null +++ b/tests/honcho_plugin/test_turn_author_peers.py @@ -0,0 +1,234 @@ +"""Tests for per-turn author attribution. + +A shared session (group, channel, thread) carries turns from several +participants and from other agents, but the manager resolves one user peer +when the session is created. Every later turn was written under that peer, +so whoever created the session collected everyone else's facts. + +``resolve_author_peer_id`` maps the turn's author onto its own peer and +``_flush_session`` writes each user message under it, joining the peer to +the Honcho session the first time it speaks. +""" + +import sys +import types +from types import SimpleNamespace +from unittest.mock import MagicMock + +import pytest + +from plugins.memory.honcho import HonchoMemoryProvider +from plugins.memory.honcho.client import HonchoClientConfig +from plugins.memory.honcho.session import HonchoSessionManager + + +def _config(**overrides) -> HonchoClientConfig: + base = dict(api_key="test-key", peer_name="eri", ai_peer="hermes") + base.update(overrides) + return HonchoClientConfig(**base) + + +def _manager(config: HonchoClientConfig, runtime_id: str | None = None) -> HonchoSessionManager: + mgr = HonchoSessionManager( + honcho=MagicMock(), + config=config, + runtime_user_peer_name=runtime_id, + ) + 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(), [])) + return mgr + + +class TestResolveAuthorPeerId: + def test_no_author_keeps_the_session_peer(self): + """An unnamed author is not attributable — never guess a peer for it.""" + mgr = _manager(_config(), runtime_id="7654321") + assert mgr.resolve_author_peer_id("telegram:group1", None) is None + assert mgr.resolve_author_peer_id("telegram:group1", "") is None + + def test_author_is_the_session_peer(self): + """The session's own participant needs no second peer.""" + mgr = _manager(_config(), runtime_id="7654321") + assert mgr.resolve_author_peer_id("telegram:group1", "7654321") is None + + def test_other_participant_gets_its_own_peer(self): + mgr = _manager(_config(), runtime_id="7654321") + assert mgr.resolve_author_peer_id("telegram:group1", "111222") == "111222" + + def test_alias_wins(self): + """An aliased account lands on its named peer whichever turn it wrote.""" + mgr = _manager( + _config(user_peer_aliases={"111222": "alice"}), + runtime_id="7654321", + ) + assert mgr.resolve_author_peer_id("telegram:group1", "111222") == "alice" + + def test_runtime_prefix_applies(self): + mgr = _manager(_config(runtime_peer_prefix="telegram_"), runtime_id="7654321") + assert mgr.resolve_author_peer_id("telegram:group1", "111222") == "telegram_111222" + + def test_pin_peer_name_collapses_authors(self): + """pinPeerName is an explicit request to unify identities.""" + mgr = _manager(_config(pin_peer_name=True), runtime_id="7654321") + assert mgr.resolve_author_peer_id("telegram:group1", "111222") is None + + def test_display_name_never_becomes_a_peer_id(self): + """Display names are attacker-influenceable on most platforms.""" + mgr = _manager(_config(), runtime_id="7654321") + assert mgr.resolve_author_peer_id("telegram:group1", None, "Alice") is None + + +class TestFlushAttributesMessages: + @pytest.fixture(autouse=True) + def _fake_sdk_session_module(self, monkeypatch): + """The join imports SessionPeerConfig from the SDK at call time and skips silently without it.""" + module = types.ModuleType("honcho.session") + module.SessionPeerConfig = lambda **kwargs: SimpleNamespace(**kwargs) + monkeypatch.setitem(sys.modules, "honcho.session", module) + + def _session(self, mgr, key="telegram:group1"): + return mgr.get_or_create(key) + + def test_author_message_written_under_the_author_peer(self): + mgr = _manager(_config(), runtime_id="7654321") + session = self._session(mgr) + honcho_session = MagicMock() + mgr._sessions_cache[session.honcho_session_id] = honcho_session + + session.add_message("user", "alice speaking", author_peer_id="alice") + assert mgr._flush_session(session) is True + + written = honcho_session.add_messages.call_args[0][0] + assert len(written) == 1 + # The peer object the message was built from is the author's, not the + # session's — that is the whole point of the change. + assert mgr._get_or_create_peer.call_args_list[-1][0][0] == "alice" + + def test_unattributed_message_keeps_the_session_peer(self): + mgr = _manager(_config(), runtime_id="7654321") + session = self._session(mgr) + honcho_session = MagicMock() + mgr._sessions_cache[session.honcho_session_id] = honcho_session + + session.add_message("user", "owner speaking") + assert mgr._flush_session(session) is True + honcho_session.add_peers.assert_not_called() + + def test_assistant_message_ignores_author(self): + """The reply is the agent's however the turn arrived.""" + mgr = _manager(_config(), runtime_id="7654321") + session = self._session(mgr) + honcho_session = MagicMock() + mgr._sessions_cache[session.honcho_session_id] = honcho_session + + session.add_message("assistant", "reply", author_peer_id="alice") + assert mgr._flush_session(session) is True + honcho_session.add_peers.assert_not_called() + + def test_author_peer_joins_once(self): + """A shared session's roster is open, so peers join when they write.""" + mgr = _manager(_config(), runtime_id="7654321") + session = self._session(mgr) + honcho_session = MagicMock() + mgr._sessions_cache[session.honcho_session_id] = honcho_session + + session.add_message("user", "first", author_peer_id="alice") + mgr._flush_session(session) + session.add_message("user", "second", author_peer_id="alice") + mgr._flush_session(session) + + assert honcho_session.add_peers.call_count == 1 + + def test_two_authors_each_join(self): + mgr = _manager(_config(), runtime_id="7654321") + session = self._session(mgr) + honcho_session = MagicMock() + mgr._sessions_cache[session.honcho_session_id] = honcho_session + + session.add_message("user", "from alice", author_peer_id="alice") + session.add_message("user", "from bob", author_peer_id="bob") + mgr._flush_session(session) + + assert honcho_session.add_peers.call_count == 2 + + def test_join_failure_still_writes_under_the_author(self): + """A failed join loses the observe config, never the attribution.""" + mgr = _manager(_config(), runtime_id="7654321") + session = self._session(mgr) + honcho_session = MagicMock() + honcho_session.add_peers.side_effect = RuntimeError("network") + mgr._sessions_cache[session.honcho_session_id] = honcho_session + + session.add_message("user", "alice speaking", author_peer_id="alice") + assert mgr._flush_session(session) is True + assert mgr._get_or_create_peer.call_args_list[-1][0][0] == "alice" + # Not remembered as joined, so the next write retries the join. + assert "alice" not in mgr._joined_author_peers.get(session.honcho_session_id, set()) + + +class TestProviderReadsTheAuthor: + def _provider(self) -> HonchoMemoryProvider: + provider = HonchoMemoryProvider() + provider._session_key = "telegram:group1" + provider._manager = MagicMock() + provider._cron_skipped = False + provider._config = SimpleNamespace(message_max_chars=25000) + return provider + + def test_on_turn_start_records_the_author(self): + provider = self._provider() + provider.on_turn_start( + 3, "hello", author_id="111222", author_name="Alice", author_is_bot=False + ) + assert provider._turn_author == { + "id": "111222", + "name": "Alice", + "is_bot": False, + } + + def test_on_turn_start_without_author_kwargs(self): + """Callers that never adopted the kwargs must keep working.""" + provider = self._provider() + provider.on_turn_start(1, "hello") + assert provider._turn_author == {"id": None, "name": None, "is_bot": False} + + def test_bot_authored_turn_is_flagged(self): + provider = self._provider() + provider.on_turn_start(2, "ping", author_id="bot-9", author_is_bot=True) + assert provider._turn_author["is_bot"] is True + + def test_sync_turn_attaches_the_resolved_author_peer(self): + provider = self._provider() + provider._session_initialized = True + session = MagicMock() + provider._manager.get_or_create.return_value = session + provider._manager.resolve_author_peer_id.return_value = "alice" + + provider.on_turn_start(1, "hi", author_id="111222", author_name="Alice") + provider.sync_turn("hi", "hello back") + if provider._sync_thread: + provider._sync_thread.join(timeout=5) + + user_calls = [ + c for c in session.add_message.call_args_list if c[0][0] == "user" + ] + assert user_calls, "the user turn was never written" + assert all(c[1]["author_peer_id"] == "alice" for c in user_calls) + + def test_sync_turn_resolves_before_the_write_thread_starts(self): + """A following turn must not retag a write that is already queued.""" + provider = self._provider() + provider._session_initialized = True + session = MagicMock() + provider._manager.get_or_create.return_value = session + provider._manager.resolve_author_peer_id.return_value = "alice" + + provider.on_turn_start(1, "hi", author_id="111222") + provider.sync_turn("hi", "hello back") + provider._turn_author = {"id": "999", "name": None, "is_bot": False} + if provider._sync_thread: + provider._sync_thread.join(timeout=5) + + provider._manager.resolve_author_peer_id.assert_called_once_with( + "telegram:group1", "111222", None + ) diff --git a/tests/test_honcho_startup_fail_open.py b/tests/test_honcho_startup_fail_open.py index 4b8e1da2d3..c4ca1d5d07 100644 --- a/tests/test_honcho_startup_fail_open.py +++ b/tests/test_honcho_startup_fail_open.py @@ -437,6 +437,9 @@ def test_honcho_sync_turn_does_not_suppress_genuine_user_messages(): manager_calls.append(session_key) return SimpleNamespace() + def resolve_author_peer_id(self, session_key, author_id, author_name=None): + return None + provider._config = _configured_tools_config(init_on_session_start=True) provider._manager = Manager() provider._session_key = "test-session" @@ -481,7 +484,10 @@ def test_honcho_sync_turn_same_instance_config_flip_gates_writes(): class Manager: def get_or_create(self, session_key): manager_calls.append(session_key) - return SimpleNamespace(add_message=lambda role, content: None) + return SimpleNamespace(add_message=lambda role, content, **kwargs: None) + + def resolve_author_peer_id(self, session_key, author_id, author_name=None): + return None def save(self, session): write_done.set()