From dc2a296bad6fbc585d670bbebe47c10ac5bd081f Mon Sep 17 00:00:00 2001 From: Padreug Date: Sat, 26 Sep 2026 23:45:47 +0200 Subject: [PATCH] fix(nostr): retry a dequeued req instead of dropping it `run_forever` took a req off the queue and then sent it; if the send raised, the req was already gone and the publish was lost outright, with the caller long since told it succeeded (the queue put returns immediately, and the publisher logs "Published" straight after). Hold the req across the reconnect and retry, bounded at three attempts so one unsendable message can't wedge every later publish behind it. This narrows but does not close the gap: a send is still confirmed at the queue, not by the relay's OK, so a half-dead socket can accept bytes that never arrive. Closing that needs OK handling in publish_nostr_event. Refs #35 --- nostr/nostr_client.py | 25 ++++++++++++++++++++++++- 1 file changed, 24 insertions(+), 1 deletion(-) diff --git a/nostr/nostr_client.py b/nostr/nostr_client.py index 4de332f..e0ae70f 100644 --- a/nostr/nostr_client.py +++ b/nostr/nostr_client.py @@ -19,6 +19,9 @@ from websocket import WebSocketApp from .event import NostrEvent 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 class NostrClient: @@ -78,17 +81,37 @@ class NostrClient: async def run_forever(self): self.running = True + # A req that was dequeued but whose send raised. It is already + # off the queue, so dropping it loses the publish outright and + # the caller has long since been told it succeeded (the queue + # put returns immediately). Hold it across the reconnect and + # retry instead. + held_req: list | None = None + held_attempts = 0 while self.running: try: if not self.is_websocket_connected: self.ws = await self.connect() await asyncio.sleep(5) - req = await self.send_req_queue.get() + if held_req is not None: + req = held_req + else: + req = await self.send_req_queue.get() + held_req, held_attempts = req, held_attempts + 1 assert self.ws self.ws.send(json.dumps(req)) + held_req, held_attempts = None, 0 except Exception as ex: logger.warning(f"[EVENTS] NostrClient error: {ex}") + if held_req is not None and held_attempts >= MAX_SEND_ATTEMPTS: + # Bounded: a req the relay or the socket will never + # accept must not wedge the queue behind it forever. + logger.error( + f"[EVENTS] Dropping req after {held_attempts} " + f"failed sends: {held_req[0]}" + ) + held_req, held_attempts = None, 0 await asyncio.sleep(60) def is_duplicate_event(self, event_id: str) -> bool: