diff --git a/plugins/platforms/buzz/adapter.py b/plugins/platforms/buzz/adapter.py index cfc8ee1a73..be90106d6a 100644 --- a/plugins/platforms/buzz/adapter.py +++ b/plugins/platforms/buzz/adapter.py @@ -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"(?`` — 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, diff --git a/tests/gateway/test_buzz_mention_resolution.py b/tests/gateway/test_buzz_mention_resolution.py new file mode 100644 index 0000000000..ba7f0de96f --- /dev/null +++ b/tests/gateway/test_buzz_mention_resolution.py @@ -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:]