fix(nostr): retry a dequeued req instead of dropping it
Some checks failed
lint.yml / fix(nostr): retry a dequeued req instead of dropping it (pull_request) Failing after 0s

`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
This commit is contained in:
Padreug 2026-09-26 23:45:47 +02:00
commit dc2a296bad

View file

@ -19,6 +19,9 @@ from websocket import WebSocketApp
from .event import NostrEvent from .event import NostrEvent
MAX_SEEN_EVENTS = 500 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: class NostrClient:
@ -78,17 +81,37 @@ class NostrClient:
async def run_forever(self): async def run_forever(self):
self.running = True 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: while self.running:
try: try:
if not self.is_websocket_connected: if not self.is_websocket_connected:
self.ws = await self.connect() self.ws = await self.connect()
await asyncio.sleep(5) await asyncio.sleep(5)
if held_req is not None:
req = held_req
else:
req = await self.send_req_queue.get() req = await self.send_req_queue.get()
held_req, held_attempts = req, held_attempts + 1
assert self.ws assert self.ws
self.ws.send(json.dumps(req)) self.ws.send(json.dumps(req))
held_req, held_attempts = None, 0
except Exception as ex: except Exception as ex:
logger.warning(f"[EVENTS] NostrClient error: {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) await asyncio.sleep(60)
def is_duplicate_event(self, event_id: str) -> bool: def is_duplicate_event(self, event_id: str) -> bool: