fix: deliver events to every matching subscription on a connection #6
2 changed files with 95 additions and 5 deletions
fix: deliver events to every matching subscription on a connection
notify_event returned after the first filter that matched, so a connection holding several subscriptions only ever received an event on one of them. That is invisible with one subscription per client, but a multiplexer such as nostrclient funnels all of its clients through a single connection. With nwcprovider subscribed to its own kind-23195 responses, every NWC reply was handed to that subscription and stopped there; the wallet app's subscription on the same connection never saw it, and Amethyst reported "wallet request timed out". Direct to the relay it worked, because the app then had its own connection. Deliver once per subscription id instead, and demote the per-filter miss log to debug: it emitted one INFO line per filter per event. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_013Tbyw6FwjhEJg3gHfPHxWt
commit
4fec1dc826
|
|
@ -72,16 +72,24 @@ class NostrClientConnection:
|
||||||
if self._is_private_event_for_other(event):
|
if self._is_private_event_for_other(event):
|
||||||
return False
|
return False
|
||||||
|
|
||||||
|
# One connection can hold several subscriptions whose filters all match
|
||||||
|
# the same event: a multiplexer such as nostrclient funnels every one of
|
||||||
|
# its clients through a single connection. NIP-01 requires the event to
|
||||||
|
# reach each matching subscription, so deliver once per subscription id
|
||||||
|
# rather than stopping at the first hit.
|
||||||
|
notified: set[str] = set()
|
||||||
for nostr_filter in self.filters:
|
for nostr_filter in self.filters:
|
||||||
|
sub_id = nostr_filter.subscription_id
|
||||||
|
if sub_id in notified:
|
||||||
|
continue
|
||||||
if nostr_filter.matches(event):
|
if nostr_filter.matches(event):
|
||||||
resp = event.serialize_response(nostr_filter.subscription_id)
|
notified.add(sub_id) # type: ignore[arg-type]
|
||||||
await self._send_msg(resp)
|
await self._send_msg(event.serialize_response(sub_id))
|
||||||
return True
|
|
||||||
else:
|
else:
|
||||||
logger.info(
|
logger.debug(
|
||||||
f"[NOSTRRELAY CLIENT] ❌ Filter didn't match for event {event.id}"
|
f"[NOSTRRELAY CLIENT] ❌ Filter didn't match for event {event.id}"
|
||||||
)
|
)
|
||||||
return False
|
return len(notified) > 0
|
||||||
|
|
||||||
def _is_private_event_for_other(self, event: NostrEvent) -> bool:
|
def _is_private_event_for_other(self, event: NostrEvent) -> bool:
|
||||||
"""
|
"""
|
||||||
|
|
|
||||||
82
tests/test_notify.py
Normal file
82
tests/test_notify.py
Normal file
|
|
@ -0,0 +1,82 @@
|
||||||
|
from unittest.mock import AsyncMock, MagicMock
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from ..relay.client_connection import NostrClientConnection
|
||||||
|
from ..relay.event import NostrEvent
|
||||||
|
from ..relay.filter import NostrFilter
|
||||||
|
from ..relay.relay import RelaySpec
|
||||||
|
|
||||||
|
RELAY_ID = "relay_notify"
|
||||||
|
WALLET = "1111111111111111111111111111111111111111111111111111111111111111"
|
||||||
|
CLIENT = "2222222222222222222222222222222222222222222222222222222222222222"
|
||||||
|
REQUEST_ID = "3" * 64
|
||||||
|
SIG = "0" * 128
|
||||||
|
|
||||||
|
|
||||||
|
def _connection() -> NostrClientConnection:
|
||||||
|
conn = NostrClientConnection(relay_id=RELAY_ID, websocket=MagicMock())
|
||||||
|
conn.get_client_config = lambda: RelaySpec()
|
||||||
|
conn._send_msg = AsyncMock() # type: ignore[method-assign]
|
||||||
|
return conn
|
||||||
|
|
||||||
|
|
||||||
|
def _nwc_response() -> NostrEvent:
|
||||||
|
return NostrEvent(
|
||||||
|
id="4" * 64,
|
||||||
|
relay_id=RELAY_ID,
|
||||||
|
publisher=WALLET,
|
||||||
|
pubkey=WALLET,
|
||||||
|
created_at=0,
|
||||||
|
kind=23195,
|
||||||
|
tags=[["p", CLIENT], ["e", REQUEST_ID]],
|
||||||
|
content="ciphertext",
|
||||||
|
sig=SIG,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _sent_subscription_ids(conn: NostrClientConnection) -> list[str]:
|
||||||
|
return [call.args[0][1] for call in conn._send_msg.await_args_list] # type: ignore[attr-defined]
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_event_reaches_every_matching_subscription():
|
||||||
|
"""
|
||||||
|
Regression: delivery used to stop at the first matching filter, so when a
|
||||||
|
wallet service (subscribed to its own kind-23195 responses) and a wallet
|
||||||
|
app shared one connection via nostrclient, the app never got the reply.
|
||||||
|
"""
|
||||||
|
conn = _connection()
|
||||||
|
conn.filters = [
|
||||||
|
NostrFilter(subscription_id="wallet-own", kinds=[23195], authors=[WALLET]),
|
||||||
|
NostrFilter(subscription_id="unrelated", kinds=[1]),
|
||||||
|
NostrFilter(
|
||||||
|
subscription_id="app-reply",
|
||||||
|
kinds=[23195],
|
||||||
|
**{"#e": [REQUEST_ID], "#p": [CLIENT]}, # type: ignore[arg-type]
|
||||||
|
),
|
||||||
|
]
|
||||||
|
|
||||||
|
assert await conn.notify_event(_nwc_response()) is True
|
||||||
|
assert _sent_subscription_ids(conn) == ["wallet-own", "app-reply"]
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_subscription_with_several_matching_filters_gets_event_once():
|
||||||
|
conn = _connection()
|
||||||
|
conn.filters = [
|
||||||
|
NostrFilter(subscription_id="multi", kinds=[23195]),
|
||||||
|
NostrFilter(subscription_id="multi", authors=[WALLET]),
|
||||||
|
]
|
||||||
|
|
||||||
|
assert await conn.notify_event(_nwc_response()) is True
|
||||||
|
assert _sent_subscription_ids(conn) == ["multi"]
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_no_matching_filter_sends_nothing():
|
||||||
|
conn = _connection()
|
||||||
|
conn.filters = [NostrFilter(subscription_id="other", kinds=[1])]
|
||||||
|
|
||||||
|
assert await conn.notify_event(_nwc_response()) is False
|
||||||
|
conn._send_msg.assert_not_awaited() # type: ignore[attr-defined]
|
||||||
Loading…
Add table
Add a link
Reference in a new issue