diff --git a/nostr/nostr_client.py b/nostr/nostr_client.py index 9afa0be..e0ae70f 100644 --- a/nostr/nostr_client.py +++ b/nostr/nostr_client.py @@ -22,11 +22,6 @@ MAX_SEEN_EVENTS = 500 # 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. 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: @@ -37,10 +32,6 @@ class NostrClient: self.subscription_id = "events-" + urlsafe_short_hash()[:32] self.running = False 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 def is_websocket_connected(self): @@ -133,82 +124,14 @@ class NostrClient: return False async def get_event(self): - """Get the next relay message, consuming `OK` frames on the way. + """Get next event from the receive queue.""" + value = await self.receive_event_queue.get() + if isinstance(value, ValueError): + raise value + return value - 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() - if isinstance(value, ValueError): - self._fail_pending_oks("connection closed") - raise value - if self._settle_ok(value): - continue - return value - - def _settle_ok(self, message) -> bool: - """Resolve the future for an `["OK", , , ]` 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()]) - 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 publish_nostr_event(self, e: NostrEvent): + await self.send_req_queue.put(["EVENT", e.dict()]) async def subscribe(self, filters: list[dict]): """Subscribe to events matching the given filters.""" diff --git a/nostr_publisher.py b/nostr_publisher.py index 677d96f..485acd4 100644 --- a/nostr_publisher.py +++ b/nostr_publisher.py @@ -203,16 +203,7 @@ async def publish_event_to_nostr( nostr_event.pubkey = signed["pubkey"] nostr_event.sig = signed["sig"] - 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 + await nostr_client.publish_nostr_event(nostr_event) logger.info( f"[EVENTS] Published NIP-52 {'delete' if delete else 'calendar'} " diff --git a/tests/test_publish_confirmation.py b/tests/test_publish_confirmation.py deleted file mode 100644 index 9615209..0000000 --- a/tests/test_publish_confirmation.py +++ /dev/null @@ -1,105 +0,0 @@ -"""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