fix(buzz): resolve @mentions to member pubkeys so agent-to-agent pings work

Salvaged from PR #83414 (4 commits squashed to final state) and composed
with the presentation-mention escape retry from PR #82646 already on this
branch: send() now resolves @Name tokens to channel-member pubkeys
(membership-accurate via `channels members`, TTL-cached, Unicode token
boundaries, ambiguous names stay presentation-only) and passes explicit
--mention args; recovery ladder handles membership drift, unresolvable
prose @tokens (escape retry, #78797), and a final self-mention downgrade.
This commit is contained in:
Matt Lutz
2026-08-31 06:22:06 -07:00
committed by Teknium
parent 1885a40ad3
commit a74080437a
2 changed files with 516 additions and 10 deletions
+212 -10
View File
@@ -232,6 +232,12 @@ _DEFAULT_POLL_INTERVAL = 4.0
_MIN_POLL_INTERVAL = 1.0
_CLI_TIMEOUT = 30.0
# Mention-resolution caches: member lists are cheap to refetch but hit on
# every publish containing "@", so a short TTL amortizes the CLI round-trip;
# display names change rarely, but must not survive a rename forever.
_MEMBER_CACHE_TTL = 60.0
_PROFILE_NAME_TTL = 300.0
# WebSocket transport (NIP-42 authenticated Nostr subscription).
# kind 44100 is Buzz's channel-membership event — used for live DM discovery.
_WS_AUTH_TIMEOUT = 20.0
@@ -894,18 +900,213 @@ class BuzzAdapter(BasePlatformAdapter):
# ── Sending ───────────────────────────────────────────────────────────
async def _run_message_send(self, args: List[str], content: str):
"""Run one send with a single unresolved-mention preflight retry."""
code, out, err = await self._run_cli(args, input_text=content)
async def _channel_member_pubkeys(self, chat_id: str) -> List[str]:
"""Candidate pubkeys for mention resolution, membership-accurate.
Primary source is ``channels members`` — the relay's membership
contract — because a ``--mention`` for a non-member makes the CLI
reject the whole publish. CLIs without that subcommand fall back
to harvesting recent channel traffic (authors plus prior mention
tags), which can over-approximate; ``send()`` recovers from any
resulting non-member mention by retrying without mention flags.
The member list is cached per channel for ``_MEMBER_CACHE_TTL``
seconds so a chatty agent doesn't pay a CLI round-trip on every
publish; membership drift inside the TTL window is covered by the
same ``send()`` recovery retry.
"""
cache: Dict[str, Tuple[float, List[str]]] = getattr(
self, "_member_cache", {}
)
self._member_cache = cache
cached = cache.get(str(chat_id))
if cached is not None and (time.monotonic() - cached[0]) < _MEMBER_CACHE_TTL:
return list(cached[1])
code, out, _err = await self._run_cli(
["channels", "members", "--channel", str(chat_id)]
)
if code == 0:
pks: List[str] = []
try:
rows = json.loads(out or "[]")
except ValueError:
rows = []
for row in rows:
pk = row.get("pubkey") if isinstance(row, dict) else row
pk = str(pk or "").lower()
if pk and pk not in pks:
pks.append(pk)
if pks:
cache[str(chat_id)] = (time.monotonic(), list(pks))
return pks
candidates: List[str] = []
code, out, _err = await self._run_cli(
["messages", "get", "--channel", str(chat_id), "--limit", "50"]
)
if code == 0:
try:
for msg in json.loads(out or "[]"):
pk = str(msg.get("pubkey") or "").lower()
if pk and pk not in candidates:
candidates.append(pk)
for t in msg.get("tags") or []:
if isinstance(t, list) and len(t) > 1 and t[0] == "p":
tpk = str(t[1]).lower()
if tpk and tpk not in candidates:
candidates.append(tpk)
except ValueError:
pass
if candidates:
cache[str(chat_id)] = (time.monotonic(), list(candidates))
return candidates
async def _profile_display_name(self, pubkey: str) -> str:
"""Display name for *pubkey* via ``users get --pubkey``, cached.
Bare ``users get`` may return only our own profile
(relay-dependent), so lookups are per-pubkey. Entries expire after
``_PROFILE_NAME_TTL`` seconds so a renamed member resolves under
their new display name without a process restart.
"""
cache: Dict[str, Tuple[float, str]] = getattr(
self, "_profile_name_cache", {}
)
self._profile_name_cache = cache
cached = cache.get(pubkey)
if cached is not None and (time.monotonic() - cached[0]) < _PROFILE_NAME_TTL:
return cached[1]
name = ""
code, out, _err = await self._run_cli(["users", "get", "--pubkey", pubkey])
if code == 0:
try:
profiles = json.loads(out or "[]")
except ValueError:
profiles = []
if profiles and isinstance(profiles[0], dict):
p0 = profiles[0]
name = str(p0.get("display_name") or p0.get("name") or "").strip()
if not name and p0.get("content"):
try:
prof = json.loads(p0["content"])
name = str(
prof.get("display_name") or prof.get("name") or ""
).strip()
except ValueError:
pass
cache[pubkey] = (time.monotonic(), name)
return name
async def _mention_pubkeys_for(self, chat_id: str, content: str) -> List[str]:
"""Resolve ``@Name`` references in *content* to member pubkeys.
The CLI hard-fails a publish when any @token fails to resolve to a
current member, and LLM prose is full of @-shaped tokens — including
real mentions with trailing punctuation ("@Riley!!") the CLI's own
parser rejects. Passing explicit ``--mention`` pubkeys for every
member name we find keeps genuine mentions notifying (p-tags intact)
while downgrading everything unresolvable to presentation-only text.
Matching is mention-token semantics, not substring, bounded on both
sides with Unicode-aware word classes: the ``@`` must start a token
("email@Fizz", "x@Fizz", "@@Fizz", and "山田@Fizz" do NOT wake
Fizz) and the name must be followed by a non-word character or
end-of-text ("@Riley!!" tags Riley; "@FizzBuzz" does NOT tag a
member named Fizz). Longer names match first and consume their
span, so "@Hermes Matt" prefers the member "Hermes Matt" over a
member "Hermes".
Duplicate display names are ambiguous: the span is consumed but no
one is tagged (presentation-only), mirroring how Buzz treats
ambiguous names — never pick an arbitrary member.
"""
if "@" not in content:
return []
by_name: Dict[str, List[str]] = {}
display: Dict[str, str] = {}
self_pk = getattr(self, "_self_pubkey", None)
for pk in await self._channel_member_pubkeys(chat_id):
if pk == self_pk:
continue
name = await self._profile_display_name(pk)
if not name:
continue
key = name.lower()
by_name.setdefault(key, [])
if pk not in by_name[key]:
by_name[key].append(pk)
display.setdefault(key, name)
found: List[str] = []
text = content
for key in sorted(by_name, key=len, reverse=True):
pattern = re.compile(
r"(?<![\w@])@" + re.escape(display[key]) + r"(?!\w)",
re.IGNORECASE,
)
if pattern.search(text):
pks = by_name[key]
if len(pks) == 1 and pks[0] not in found:
found.append(pks[0])
# Consume the span either way: a shorter member name that is
# a prefix of this one must not double-match, and an
# ambiguous name must stay presentation-only rather than
# falling through to a partial match.
text = pattern.sub("\x00", text)
return found
async def _run_message_send(
self,
args: List[str],
content: str,
mention_pubkeys: Optional[List[str]] = None,
):
"""Run one send with bounded mention-failure recovery.
Ladder (each rung fires at most once):
1. publish with explicit ``--mention`` pubkeys resolved from the
content (#83414) so genuine member mentions carry p-tags and
mention-subscribed agents actually wake;
2. if the CLI rejects because a resolved pubkey is no longer a
member (membership drift), retry without the explicit mentions —
deliver the message rather than lose it;
3. if the CLI's preflight rejects an unresolvable presentation
``@token`` in prose, escape exactly that token with an invisible
separator and retry (#82646 / #78797);
4. if the error persists and we know our own pubkey, retry once with
``--mention <self>`` — supplying any explicit identity downgrades
unresolvable @names to presentation-only text (#83414); the echo
de-dupe already suppresses self-notification.
"""
mention_args: List[str] = []
for pk in mention_pubkeys or []:
mention_args += ["--mention", pk]
code, out, err = await self._run_cli(args + mention_args, input_text=content)
if code == 0:
return code, out, err
if mention_args and "not channel members" in (err or ""):
# Membership drifted between resolution and publish (or the
# fallback candidate source over-approximated): never let a
# stale mention kill the message.
code, out, err = await self._run_cli(args, input_text=content)
if code == 0:
return code, out, err
escaped = _escape_unresolved_presentation_mention(content, err)
if escaped is None:
return code, out, err
logger.info(
"Buzz: retrying message after unresolved presentation-mention preflight"
)
return await self._run_cli(args, input_text=escaped)
if escaped is not None:
logger.info(
"Buzz: retrying message after unresolved presentation-mention preflight"
)
code, out, err = await self._run_cli(args, input_text=escaped)
if code == 0:
return code, out, err
if (
code != 0
and "does not match a current channel member" in (err or "")
and getattr(self, "_self_pubkey", None)
):
code, out, err = await self._run_cli(
args + ["--mention", self._self_pubkey], input_text=content
)
return code, out, err
async def send(
self,
@@ -927,7 +1128,8 @@ class BuzzAdapter(BasePlatformAdapter):
)
if reply_target and self._reply_to_mode != "off":
args += ["--reply-to", str(reply_target)]
code, out, err = await self._run_message_send(args, content)
mention_pubkeys = await self._mention_pubkeys_for(chat_id, content)
code, out, err = await self._run_message_send(args, content, mention_pubkeys)
if code != 0:
return SendResult(
success=False,
@@ -0,0 +1,304 @@
"""Tests for the Buzz adapter's @mention resolution (PR #83414).
Covers ``_mention_pubkeys_for`` / ``_channel_member_pubkeys`` and the
``send()`` recovery paths: membership-accurate candidate sourcing, token
boundaries on both sides of the @name, duplicate-name ambiguity, and the
non-member / unresolvable-token publish retries.
"""
import json
import pytest
from tests.gateway._plugin_adapter_loader import load_plugin_adapter
_buzz_mod = load_plugin_adapter("buzz")
BuzzAdapter = _buzz_mod.BuzzAdapter
SELF_PUBKEY = "9fd5c7ba6d3ef224da78f541e0fcb9c50f72cc63edb19aae76ac6a0474dfa860"
FIZZ_PUBKEY = "b" * 64
BUZZ_PUBKEY = "c" * 64
DUPE_PUBKEY = "d" * 64
CHANNEL = "ccc2bc1a-7a82-5a8f-8c4e-57a070cbe7cd"
_ENV_VARS = (
"BUZZ_RELAY_URL",
"BUZZ_PRIVATE_KEY",
"BUZZ_CHANNELS",
"BUZZ_HOME_CHANNEL",
"BUZZ_ALLOWED_USERS",
"BUZZ_ALLOW_ALL_USERS",
"BUZZ_POLL_INTERVAL",
"BUZZ_CLI_PATH",
"BUZZ_CREDENTIALS_FILE",
)
@pytest.fixture(autouse=True)
def _clean_env(monkeypatch, tmp_path):
for var in _ENV_VARS:
monkeypatch.delenv(var, raising=False)
monkeypatch.setattr(_buzz_mod, "_DEFAULT_CREDENTIALS_DIR", tmp_path / "no-creds")
yield
def _make_adapter(extra=None):
from gateway.config import PlatformConfig
cfg = PlatformConfig(enabled=True, extra={"relay_url": "https://test.relay", **(extra or {})})
adapter = BuzzAdapter(cfg)
adapter._self_pubkey = SELF_PUBKEY
adapter._private_key = "nsec1test"
return adapter
class _ScriptedCli:
"""Fake ``_run_cli`` that routes on the buzz subcommand and records calls."""
def __init__(self):
self.responses = {}
self.calls = []
def script(self, group, cmd, payload, code=0, stderr=""):
stdout = payload if isinstance(payload, str) else json.dumps(payload)
self.responses.setdefault((group, cmd), []).append((code, stdout, stderr))
async def __call__(self, args, *, input_text=None):
self.calls.append((list(args), input_text))
queue = self.responses.get((args[0], args[1]), [])
if len(queue) > 1:
return queue.pop(0)
if queue:
return queue[0]
return 0, "[]", ""
def _wire(adapter, cli):
adapter._run_cli = cli
return cli
def _members(*pubkeys):
return [{"pubkey": pk} for pk in pubkeys]
def _profile(pubkey, name):
return [{"pubkey": pubkey, "display_name": name}]
# ── candidate sourcing ────────────────────────────────────────────────────
class TestChannelMemberPubkeys:
@pytest.mark.asyncio
async def test_members_subcommand_is_primary_source(self):
adapter = _make_adapter()
cli = _wire(adapter, _ScriptedCli())
cli.script("channels", "members", _members(FIZZ_PUBKEY, BUZZ_PUBKEY))
pks = await adapter._channel_member_pubkeys(CHANNEL)
assert pks == [FIZZ_PUBKEY, BUZZ_PUBKEY]
assert ("messages", "get") not in {(c[0][0], c[0][1]) for c in cli.calls}
@pytest.mark.asyncio
async def test_falls_back_to_recent_traffic_when_members_unavailable(self):
adapter = _make_adapter()
cli = _wire(adapter, _ScriptedCli())
cli.script("channels", "members", "", code=1, stderr="unknown subcommand")
cli.script(
"messages",
"get",
[
{"pubkey": FIZZ_PUBKEY, "tags": [["p", BUZZ_PUBKEY]]},
],
)
pks = await adapter._channel_member_pubkeys(CHANNEL)
assert FIZZ_PUBKEY in pks and BUZZ_PUBKEY in pks
# ── token matching ────────────────────────────────────────────────────────
class TestMentionTokenMatching:
async def _resolve(self, content, members=None, profiles=None):
adapter = _make_adapter()
cli = _wire(adapter, _ScriptedCli())
cli.script("channels", "members", _members(*(members or [FIZZ_PUBKEY])))
for pk, name in (profiles or {FIZZ_PUBKEY: "Fizz"}).items():
cli.script("users", "get", _profile(pk, name))
return await adapter._mention_pubkeys_for(CHANNEL, content)
@pytest.mark.asyncio
async def test_clean_mention_resolves(self):
assert await self._resolve("hey @Fizz, ping") == [FIZZ_PUBKEY]
@pytest.mark.asyncio
async def test_trailing_punctuation_resolves(self):
assert await self._resolve("@Fizz!! wake up") == [FIZZ_PUBKEY]
@pytest.mark.asyncio
async def test_right_boundary_rejects_longer_token(self):
assert await self._resolve("@FizzBuzz is a game") == []
@pytest.mark.asyncio
async def test_left_boundary_rejects_email_like_text(self):
assert await self._resolve("mail me at email@Fizz today") == []
@pytest.mark.asyncio
async def test_left_boundary_rejects_adjacent_word_char(self):
assert await self._resolve("x@Fizz") == []
@pytest.mark.asyncio
async def test_left_boundary_rejects_double_at(self):
assert await self._resolve("@@Fizz") == []
@pytest.mark.asyncio
async def test_left_boundary_is_unicode_aware(self):
# 山 is a word character: "山田@Fizz" is email-shaped text in a
# non-ASCII script, not a mention.
assert await self._resolve("山田@Fizz") == []
@pytest.mark.asyncio
async def test_no_at_sign_short_circuits_without_cli_calls(self):
adapter = _make_adapter()
cli = _wire(adapter, _ScriptedCli())
assert await adapter._mention_pubkeys_for(CHANNEL, "no mentions here") == []
assert cli.calls == []
@pytest.mark.asyncio
async def test_longest_name_wins_and_consumes_span(self):
result = await self._resolve(
"@Hermes Matt please review",
members=[FIZZ_PUBKEY, BUZZ_PUBKEY],
profiles={FIZZ_PUBKEY: "Hermes Matt", BUZZ_PUBKEY: "Hermes"},
)
assert result == [FIZZ_PUBKEY]
@pytest.mark.asyncio
async def test_duplicate_display_names_tag_nobody(self):
result = await self._resolve(
"@Fizz which one of you is real",
members=[FIZZ_PUBKEY, DUPE_PUBKEY],
profiles={FIZZ_PUBKEY: "Fizz", DUPE_PUBKEY: "Fizz"},
)
assert result == []
@pytest.mark.asyncio
async def test_self_is_never_a_candidate(self):
result = await self._resolve(
"@Chip and @Fizz",
members=[SELF_PUBKEY, FIZZ_PUBKEY],
profiles={SELF_PUBKEY: "Chip", FIZZ_PUBKEY: "Fizz"},
)
assert result == [FIZZ_PUBKEY]
# ── caching ───────────────────────────────────────────────────────────────
class TestResolutionCaching:
@pytest.mark.asyncio
async def test_member_list_cached_within_ttl(self):
adapter = _make_adapter()
cli = _wire(adapter, _ScriptedCli())
cli.script("channels", "members", _members(FIZZ_PUBKEY))
cli.script("users", "get", _profile(FIZZ_PUBKEY, "Fizz"))
assert await adapter._mention_pubkeys_for(CHANNEL, "@Fizz one") == [FIZZ_PUBKEY]
assert await adapter._mention_pubkeys_for(CHANNEL, "@Fizz two") == [FIZZ_PUBKEY]
member_calls = [c for c in cli.calls if (c[0][0], c[0][1]) == ("channels", "members")]
assert len(member_calls) == 1, "second resolve must hit the member cache"
@pytest.mark.asyncio
async def test_profile_name_expires_after_ttl(self, monkeypatch):
adapter = _make_adapter()
cli = _wire(adapter, _ScriptedCli())
cli.script("channels", "members", _members(FIZZ_PUBKEY))
cli.script("users", "get", _profile(FIZZ_PUBKEY, "Fizz"))
cli.script("users", "get", _profile(FIZZ_PUBKEY, "FizzRenamed"))
clock = [1000.0]
monkeypatch.setattr(_buzz_mod.time, "monotonic", lambda: clock[0])
assert await adapter._mention_pubkeys_for(CHANNEL, "@Fizz hi") == [FIZZ_PUBKEY]
# Inside both TTLs: rename not visible yet, no new lookups needed.
assert await adapter._mention_pubkeys_for(CHANNEL, "@FizzRenamed hi") == []
# Past the name TTL (and member TTL): the rename resolves.
clock[0] += _buzz_mod._PROFILE_NAME_TTL + 1
assert await adapter._mention_pubkeys_for(CHANNEL, "@FizzRenamed hi") == [FIZZ_PUBKEY]
# ── send() recovery paths ─────────────────────────────────────────────────
class TestSendRecovery:
def _sending_adapter(self):
adapter = _make_adapter()
cli = _wire(adapter, _ScriptedCli())
cli.script("channels", "members", _members(FIZZ_PUBKEY))
cli.script("users", "get", _profile(FIZZ_PUBKEY, "Fizz"))
return adapter, cli
def _send_calls(self, cli):
return [c for c in cli.calls if (c[0][0], c[0][1]) == ("messages", "send")]
@pytest.mark.asyncio
async def test_resolved_mention_attached_to_publish(self):
adapter, cli = self._sending_adapter()
cli.script("messages", "send", {"accepted": True, "event_id": "e1"})
result = await adapter.send(CHANNEL, "@Fizz hello")
assert result.success
sends = self._send_calls(cli)
assert len(sends) == 1
assert ["--mention", FIZZ_PUBKEY] == sends[0][0][-2:]
@pytest.mark.asyncio
async def test_non_member_mention_retries_without_mentions(self):
adapter, cli = self._sending_adapter()
cli.script(
"messages", "send", "",
code=1,
stderr=json.dumps({"error": "user_error", "message": "mentioned pubkeys are not channel members"}),
)
cli.script("messages", "send", {"accepted": True, "event_id": "e2"})
result = await adapter.send(CHANNEL, "@Fizz hello")
assert result.success
sends = self._send_calls(cli)
assert len(sends) == 2
assert "--mention" in sends[0][0]
assert "--mention" not in sends[1][0]
@pytest.mark.asyncio
async def test_unresolvable_token_retries_with_self_mention(self):
adapter = _make_adapter()
cli = _wire(adapter, _ScriptedCli())
cli.script("channels", "members", _members()) # nobody to resolve
cli.script(
"messages", "send", "",
code=1,
stderr=json.dumps({
"error": "user_error",
"message": "mention '@-mention' does not match a current channel member",
}),
)
cli.script("messages", "send", {"accepted": True, "event_id": "e3"})
result = await adapter.send(CHANNEL, "just @-mention me")
assert result.success
sends = self._send_calls(cli)
assert len(sends) == 2
assert ["--mention", SELF_PUBKEY] == sends[1][0][-2:]