feat(honcho): write each turn under its author's peer
The manager resolved one user peer in `get_or_create` and froze it onto the session, then `_flush_session` chose between it and the assistant peer by role. Every user turn in a shared session landed on that one peer, so the first person to message the agent collected everyone else's facts — and a Honcho conclusion, once derived, is not self-correcting. `resolve_author_peer_id` maps the turn's author onto its own peer using the alias-then-prefix order `_resolve_user_peer_id` already applies, so an aliased account reaches the same peer whichever turn it wrote. `sync_turn` resolves it before starting the write thread, so a following turn cannot retag a queued write. `_flush_session` then writes each user message under that peer. A shared session's roster is open — people and other agents arrive after the session exists — so `_author_peer_for_session` joins a peer when it first writes instead of enumerating participants at init. Joins are remembered per session, and a failed join still writes under the right peer, losing only the observe config. Three cases return None and keep the session's own peer: no author named, the author IS the session's peer, and `pinPeerName` set — that flag is an explicit request to unify identities, so it still collapses authors in a shared chat. An unnamed author stays unattributed rather than defaulting to the owner. That preserves today's behavior for the transports that send no author, so those turns still reach the session peer; #83500 owner-gated the memory-file migration for the same reason. Display names never become peer IDs — they are attacker-influenceable on any platform where participants set their own name.
This commit is contained in:
@@ -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)
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
234
tests/honcho_plugin/test_turn_author_peers.py
Normal file
234
tests/honcho_plugin/test_turn_author_peers.py
Normal file
@@ -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
|
||||
)
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user