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: