Merge pull request 'feat(nostr): confirm publishes against the relay's OK' (#58) from feat/confirm-publish-with-relay-ok into main
Some checks failed
lint.yml / Merge pull request 'feat(nostr): confirm publishes against the relay's OK' (#58) from feat/confirm-publish-with-relay-ok into main (push) Failing after 0s
Some checks failed
lint.yml / Merge pull request 'feat(nostr): confirm publishes against the relay's OK' (#58) from feat/confirm-publish-with-relay-ok into main (push) Failing after 0s
Reviewed-on: #58
This commit is contained in:
commit
13eb549066
3 changed files with 200 additions and 9 deletions
|
|
@ -22,6 +22,11 @@ MAX_SEEN_EVENTS = 500
|
||||||
# How many times a dequeued req is retried before it is dropped. Bounded
|
# How many times a dequeued req is retried before it is dropped. Bounded
|
||||||
# so one unsendable message can't block every later publish behind it.
|
# so one unsendable message can't block every later publish behind it.
|
||||||
MAX_SEND_ATTEMPTS = 3
|
MAX_SEND_ATTEMPTS = 3
|
||||||
|
# How long to wait for the relay's `OK` before treating a publish as
|
||||||
|
# unconfirmed. nostrclient's router answers within its own
|
||||||
|
# PUBLISH_TIMEOUT_SECONDS (10s) even when every relay stays silent, so
|
||||||
|
# this only needs headroom over that.
|
||||||
|
PUBLISH_OK_TIMEOUT_SECONDS = 12
|
||||||
|
|
||||||
|
|
||||||
class NostrClient:
|
class NostrClient:
|
||||||
|
|
@ -32,6 +37,10 @@ class NostrClient:
|
||||||
self.subscription_id = "events-" + urlsafe_short_hash()[:32]
|
self.subscription_id = "events-" + urlsafe_short_hash()[:32]
|
||||||
self.running = False
|
self.running = False
|
||||||
self._seen_events: OrderedDict[str, None] = OrderedDict()
|
self._seen_events: OrderedDict[str, None] = OrderedDict()
|
||||||
|
# event id -> future awaiting that publish's `OK`. Resolved in
|
||||||
|
# `get_event`, which is the single point where relay messages
|
||||||
|
# cross into the event loop.
|
||||||
|
self._pending_oks: dict[str, asyncio.Future] = {}
|
||||||
|
|
||||||
@property
|
@property
|
||||||
def is_websocket_connected(self):
|
def is_websocket_connected(self):
|
||||||
|
|
@ -124,14 +133,82 @@ class NostrClient:
|
||||||
return False
|
return False
|
||||||
|
|
||||||
async def get_event(self):
|
async def get_event(self):
|
||||||
"""Get next event from the receive queue."""
|
"""Get the next relay message, consuming `OK` frames on the way.
|
||||||
|
|
||||||
|
This is the only place relay messages cross from the websocket
|
||||||
|
thread into the event loop, which makes it the natural place to
|
||||||
|
settle publish confirmations — no cross-thread future juggling.
|
||||||
|
`OK` frames are swallowed rather than forwarded; the sync loop
|
||||||
|
never handled them.
|
||||||
|
"""
|
||||||
|
while True:
|
||||||
value = await self.receive_event_queue.get()
|
value = await self.receive_event_queue.get()
|
||||||
if isinstance(value, ValueError):
|
if isinstance(value, ValueError):
|
||||||
|
self._fail_pending_oks("connection closed")
|
||||||
raise value
|
raise value
|
||||||
|
if self._settle_ok(value):
|
||||||
|
continue
|
||||||
return value
|
return value
|
||||||
|
|
||||||
async def publish_nostr_event(self, e: NostrEvent):
|
def _settle_ok(self, message) -> bool:
|
||||||
|
"""Resolve the future for an `["OK", <id>, <bool>, <msg>]` frame.
|
||||||
|
|
||||||
|
Returns True when `message` was an OK frame (and so should not
|
||||||
|
be forwarded), False otherwise. An OK for a publish we are not
|
||||||
|
waiting on — a retry whose original already timed out, say — is
|
||||||
|
still consumed; it has nowhere useful to go.
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
data = json.loads(message)
|
||||||
|
except (json.JSONDecodeError, TypeError):
|
||||||
|
return False
|
||||||
|
if not isinstance(data, list) or len(data) < 3 or data[0] != "OK":
|
||||||
|
return False
|
||||||
|
|
||||||
|
event_id = data[1]
|
||||||
|
accepted = bool(data[2])
|
||||||
|
detail = data[3] if len(data) > 3 and isinstance(data[3], str) else ""
|
||||||
|
future = self._pending_oks.get(event_id)
|
||||||
|
if future and not future.done():
|
||||||
|
future.set_result((accepted, detail))
|
||||||
|
return True
|
||||||
|
|
||||||
|
def _fail_pending_oks(self, reason: str) -> None:
|
||||||
|
"""Settle every in-flight publish as unconfirmed.
|
||||||
|
|
||||||
|
Without this a disconnect leaves callers waiting the full
|
||||||
|
timeout for an `OK` that can no longer arrive.
|
||||||
|
"""
|
||||||
|
for future in self._pending_oks.values():
|
||||||
|
if not future.done():
|
||||||
|
future.set_result((False, f"error: {reason}"))
|
||||||
|
|
||||||
|
async def publish_nostr_event(self, e: NostrEvent) -> bool:
|
||||||
|
"""Publish and wait for the relay's `OK`. True only when accepted.
|
||||||
|
|
||||||
|
Queueing is not delivery: nostrclient drops an EVENT outright
|
||||||
|
when no relay is connected, answering `OK false`. Reporting
|
||||||
|
success on the queue put let a stale ticket count survive a
|
||||||
|
republish that never left the building (aiolabs/events#56).
|
||||||
|
"""
|
||||||
|
future: asyncio.Future = asyncio.get_running_loop().create_future()
|
||||||
|
self._pending_oks[e.id] = future
|
||||||
|
try:
|
||||||
await self.send_req_queue.put(["EVENT", e.dict()])
|
await self.send_req_queue.put(["EVENT", e.dict()])
|
||||||
|
accepted, detail = await asyncio.wait_for(
|
||||||
|
future, PUBLISH_OK_TIMEOUT_SECONDS
|
||||||
|
)
|
||||||
|
if not accepted:
|
||||||
|
logger.warning(f"[EVENTS] Relay rejected event {e.id[:12]}…: {detail}")
|
||||||
|
return accepted
|
||||||
|
except asyncio.TimeoutError:
|
||||||
|
logger.warning(
|
||||||
|
f"[EVENTS] No OK for event {e.id[:12]}… within "
|
||||||
|
f"{PUBLISH_OK_TIMEOUT_SECONDS}s — treating as unconfirmed"
|
||||||
|
)
|
||||||
|
return False
|
||||||
|
finally:
|
||||||
|
self._pending_oks.pop(e.id, None)
|
||||||
|
|
||||||
async def subscribe(self, filters: list[dict]):
|
async def subscribe(self, filters: list[dict]):
|
||||||
"""Subscribe to events matching the given filters."""
|
"""Subscribe to events matching the given filters."""
|
||||||
|
|
|
||||||
|
|
@ -203,7 +203,16 @@ async def publish_event_to_nostr(
|
||||||
nostr_event.pubkey = signed["pubkey"]
|
nostr_event.pubkey = signed["pubkey"]
|
||||||
nostr_event.sig = signed["sig"]
|
nostr_event.sig = signed["sig"]
|
||||||
|
|
||||||
await nostr_client.publish_nostr_event(nostr_event)
|
accepted = await nostr_client.publish_nostr_event(nostr_event)
|
||||||
|
if not accepted:
|
||||||
|
# Returning None keeps `nostr_publish_pending` set, so the
|
||||||
|
# sweep retries instead of recording a delivery that never
|
||||||
|
# happened (aiolabs/events#56).
|
||||||
|
logger.warning(
|
||||||
|
f"[EVENTS] Relay did not confirm NIP-52 "
|
||||||
|
f"{'delete' if delete else 'calendar'} event for {event.id}"
|
||||||
|
)
|
||||||
|
return None
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
f"[EVENTS] Published NIP-52 {'delete' if delete else 'calendar'} "
|
f"[EVENTS] Published NIP-52 {'delete' if delete else 'calendar'} "
|
||||||
|
|
|
||||||
105
tests/test_publish_confirmation.py
Normal file
105
tests/test_publish_confirmation.py
Normal file
|
|
@ -0,0 +1,105 @@
|
||||||
|
"""Publish confirmation against the relay's `OK` (aiolabs/events#56).
|
||||||
|
|
||||||
|
Queueing is not delivery. nostrclient drops an EVENT outright when no
|
||||||
|
relay is connected and answers `OK false`; before this, that reply was
|
||||||
|
discarded and the publish reported success, which once cleared the
|
||||||
|
`nostr_publish_pending` flag on a republish that never left the
|
||||||
|
building.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import json
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from ..nostr import nostr_client as nc
|
||||||
|
from ..nostr.event import NostrEvent
|
||||||
|
|
||||||
|
|
||||||
|
def _event(event_id: str = "a" * 64) -> NostrEvent:
|
||||||
|
e = NostrEvent(pubkey="b" * 64, created_at=0, kind=31923)
|
||||||
|
e.id = event_id
|
||||||
|
return e
|
||||||
|
|
||||||
|
|
||||||
|
async def _publish_and_reply(client, event, reply, delay=0.01):
|
||||||
|
"""Start a publish, then feed `reply` through the receive path."""
|
||||||
|
task = asyncio.create_task(client.publish_nostr_event(event))
|
||||||
|
await asyncio.sleep(delay) # let the future register
|
||||||
|
if reply is not None:
|
||||||
|
client.receive_event_queue.put_nowait(reply)
|
||||||
|
consumer = asyncio.create_task(client.get_event())
|
||||||
|
await asyncio.sleep(delay)
|
||||||
|
consumer.cancel()
|
||||||
|
return await task
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_accepted_publish_returns_true():
|
||||||
|
client = nc.NostrClient()
|
||||||
|
event = _event()
|
||||||
|
ok = json.dumps(["OK", event.id, True, ""])
|
||||||
|
assert await _publish_and_reply(client, event, ok) is True
|
||||||
|
assert client._pending_oks == {}
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_rejected_publish_returns_false():
|
||||||
|
"""The shape that bit us: no relay connected, so nostrclient's
|
||||||
|
router answers `OK false` without the event ever being sent."""
|
||||||
|
client = nc.NostrClient()
|
||||||
|
event = _event()
|
||||||
|
ok = json.dumps(["OK", event.id, False, "error: no relays connected"])
|
||||||
|
assert await _publish_and_reply(client, event, ok) is False
|
||||||
|
assert client._pending_oks == {}
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_missing_ok_times_out_as_unconfirmed(monkeypatch):
|
||||||
|
monkeypatch.setattr(nc, "PUBLISH_OK_TIMEOUT_SECONDS", 0.05)
|
||||||
|
client = nc.NostrClient()
|
||||||
|
assert await _publish_and_reply(client, _event(), None) is False
|
||||||
|
assert client._pending_oks == {}
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_ok_for_a_different_event_does_not_settle_ours(monkeypatch):
|
||||||
|
monkeypatch.setattr(nc, "PUBLISH_OK_TIMEOUT_SECONDS", 0.05)
|
||||||
|
client = nc.NostrClient()
|
||||||
|
other = json.dumps(["OK", "c" * 64, True, ""])
|
||||||
|
assert await _publish_and_reply(client, _event(), other) is False
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_get_event_swallows_ok_and_forwards_everything_else():
|
||||||
|
client = nc.NostrClient()
|
||||||
|
client.receive_event_queue.put_nowait(json.dumps(["OK", "d" * 64, True, ""]))
|
||||||
|
forwarded = json.dumps(["EVENT", "sub", {"id": "e" * 64}])
|
||||||
|
client.receive_event_queue.put_nowait(forwarded)
|
||||||
|
assert await client.get_event() == forwarded
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_disconnect_settles_inflight_publishes_immediately(monkeypatch):
|
||||||
|
"""A dropped socket must not leave the caller waiting the full
|
||||||
|
timeout for an OK that can no longer arrive."""
|
||||||
|
monkeypatch.setattr(nc, "PUBLISH_OK_TIMEOUT_SECONDS", 30)
|
||||||
|
client = nc.NostrClient()
|
||||||
|
event = _event()
|
||||||
|
task = asyncio.create_task(client.publish_nostr_event(event))
|
||||||
|
await asyncio.sleep(0.01)
|
||||||
|
|
||||||
|
client.receive_event_queue.put_nowait(ValueError("WebSocket closed"))
|
||||||
|
consumer = asyncio.create_task(client.get_event())
|
||||||
|
await asyncio.sleep(0.01)
|
||||||
|
consumer.cancel()
|
||||||
|
|
||||||
|
assert await asyncio.wait_for(task, 1) is False
|
||||||
|
|
||||||
|
|
||||||
|
def test_settle_ok_ignores_non_ok_frames():
|
||||||
|
client = nc.NostrClient()
|
||||||
|
assert client._settle_ok(json.dumps(["EVENT", "sub", {}])) is False
|
||||||
|
assert client._settle_ok(json.dumps(["EOSE", "sub"])) is False
|
||||||
|
assert client._settle_ok("not json") is False
|
||||||
|
assert client._settle_ok(json.dumps(["OK", "f" * 64, True, ""])) is True
|
||||||
Loading…
Add table
Add a link
Reference in a new issue