fix: deliver events to every matching subscription on a connection #6
4 changed files with 104 additions and 6 deletions
|
|
@ -1,6 +1,6 @@
|
||||||
{
|
{
|
||||||
"name": "Nostr Relay",
|
"name": "Nostr Relay",
|
||||||
"version": "1.1.0",
|
"version": "1.1.0-aio.3",
|
||||||
"short_description": "One click launch your own relay!",
|
"short_description": "One click launch your own relay!",
|
||||||
"tile": "/nostrrelay/static/image/nostrrelay.png",
|
"tile": "/nostrrelay/static/image/nostrrelay.png",
|
||||||
"min_lnbits_version": "1.4.0",
|
"min_lnbits_version": "1.4.0",
|
||||||
|
|
|
||||||
8
docs/upstream-candidates.md
Normal file
8
docs/upstream-candidates.md
Normal file
|
|
@ -0,0 +1,8 @@
|
||||||
|
# Upstream candidates
|
||||||
|
|
||||||
|
Fork changes that `lnbits/nostrrelay` might want, kept in upstream's shape
|
||||||
|
so a future rebase merges cleanly. Strike a row when the upstream PR merges.
|
||||||
|
|
||||||
|
| Landed | Change | Fork commit | Upstream status |
|
||||||
|
| --- | --- | --- | --- |
|
||||||
|
| 2026-09-12 | `notify_event` delivers to every matching subscription on a connection instead of stopping at the first match (breaks any multiplexed client, e.g. nostrclient + nwcprovider) | see `git log -- relay/client_connection.py` | not yet proposed; upstream `a87bc1f` has the same first-match `return True` |
|
||||||
|
|
@ -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