From cc730256abe921665068a81c441145a6f797a932 Mon Sep 17 00:00:00 2001 From: Padreug Date: Sun, 27 Sep 2026 12:36:55 +0200 Subject: [PATCH] feat(nostr): confirm publishes against the relay's OK MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `publish_nostr_event` returned as soon as the EVENT was on the send queue, and the publisher logged "Published" on the next line. Queueing is not delivery: nostrclient drops an EVENT outright when no relay is connected, answering `OK false "error: no relays connected"`. We threw that reply away. On cfaun this cost a completed repair. The #55 sweep republished a 14-day-stale calendar event 21 seconds before nostrclient had finished connecting to its relay, got `OK false`, logged `Published`, reported `1/1 recovered` and cleared `nostr_publish_pending` — leaving the count stale, the row unflagged and the log asserting success. It took a manual re-arm of the flag to finish the job. So the flag's contract was never true: it claimed to clear only on a confirmed success but cleared on a confirmed enqueue. `publish_nostr_event` now registers a future per event id, awaits the `OK`, and returns whether it was accepted. `publish_event_to_nostr` returns None when unconfirmed, which keeps the row flagged so the sweep retries rather than recording a delivery that never happened. Correlation lives in `get_event`, the one place relay messages cross from the websocket thread into the event loop — no cross-thread future juggling. OK frames are consumed there rather than forwarded; the sync loop never handled them. A disconnect settles every in-flight publish immediately instead of making callers wait out the timeout. On latency: the timeout is not the common cost. A disconnected relay is rejected by nostrclient's router in milliseconds (230ms measured on aio-demo), so the 12s budget only applies when relays are connected but silent, which nostrclient itself bounds at 10s. `set_ticket_paid` runs on the invoice-listener task, so that narrow case does stall the loop; if it ever matters, the remedy is to stop awaiting on the sale path while leaving the flag set — the sweep already guarantees eventual delivery — not to go back to reporting unverified success. Closes #56 --- nostr/nostr_client.py | 91 +++++++++++++++++++++++-- nostr_publisher.py | 11 ++- tests/test_publish_confirmation.py | 105 +++++++++++++++++++++++++++++ 3 files changed, 199 insertions(+), 8 deletions(-) create mode 100644 tests/test_publish_confirmation.py diff --git a/nostr/nostr_client.py b/nostr/nostr_client.py index e0ae70f..9afa0be 100644 --- a/nostr/nostr_client.py +++ b/nostr/nostr_client.py @@ -22,6 +22,11 @@ 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: @@ -32,6 +37,10 @@ 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): @@ -124,14 +133,82 @@ class NostrClient: return False async def get_event(self): - """Get next event from the receive queue.""" - value = await self.receive_event_queue.get() - if isinstance(value, ValueError): - raise value - return value + """Get the next relay message, consuming `OK` frames on the way. - async def publish_nostr_event(self, e: NostrEvent): - await self.send_req_queue.put(["EVENT", e.dict()]) + 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 subscribe(self, filters: list[dict]): """Subscribe to events matching the given filters.""" diff --git a/nostr_publisher.py b/nostr_publisher.py index 485acd4..677d96f 100644 --- a/nostr_publisher.py +++ b/nostr_publisher.py @@ -203,7 +203,16 @@ async def publish_event_to_nostr( nostr_event.pubkey = signed["pubkey"] 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( f"[EVENTS] Published NIP-52 {'delete' if delete else 'calendar'} " diff --git a/tests/test_publish_confirmation.py b/tests/test_publish_confirmation.py new file mode 100644 index 0000000..9615209 --- /dev/null +++ b/tests/test_publish_confirmation.py @@ -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 -- 2.55.0