feat(nostr): confirm publishes against the relay's OK
Some checks failed
lint.yml / feat(nostr): confirm publishes against the relay's OK (pull_request) Failing after 0s

`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
This commit is contained in:
Padreug 2026-09-27 12:36:55 +02:00
commit cc730256ab
3 changed files with 200 additions and 9 deletions

View file

@ -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", <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()])
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."""