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
235 lines
8.7 KiB
Python
235 lines
8.7 KiB
Python
"""
|
|
Bidirectional Nostr client for the events extension.
|
|
|
|
Connects to the nostrclient extension's internal WebSocket to publish
|
|
and subscribe to NIP-52 calendar events. Based on nostrmarket's
|
|
NostrClient pattern.
|
|
"""
|
|
|
|
import asyncio
|
|
import json
|
|
from asyncio import Queue
|
|
from collections import OrderedDict
|
|
|
|
from lnbits.helpers import encrypt_internal_message, urlsafe_short_hash
|
|
from lnbits.settings import settings
|
|
from loguru import logger
|
|
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
|
|
# 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:
|
|
def __init__(self):
|
|
self.receive_event_queue: Queue = Queue()
|
|
self.send_req_queue: Queue = Queue()
|
|
self.ws: WebSocketApp | None = None
|
|
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):
|
|
if not self.ws:
|
|
return False
|
|
return self.ws.keep_running
|
|
|
|
async def connect(self) -> WebSocketApp:
|
|
relay_endpoint = encrypt_internal_message("relay", urlsafe=True)
|
|
ws_url = (
|
|
f"ws://localhost:{settings.port}" f"/nostrclient/api/v1/{relay_endpoint}"
|
|
)
|
|
|
|
logger.info("[EVENTS] Connecting to nostrclient WebSocket...")
|
|
|
|
def on_open(_):
|
|
logger.info("[EVENTS] Connected to nostrclient WebSocket")
|
|
|
|
def on_message(_, message):
|
|
try:
|
|
self.receive_event_queue.put_nowait(message)
|
|
except Exception as e:
|
|
logger.error(f"[EVENTS] Failed to queue message: {e}")
|
|
|
|
def on_error(_, error):
|
|
logger.warning(f"[EVENTS] WebSocket error: {error}")
|
|
|
|
def on_close(_, status_code, message):
|
|
logger.warning(f"[EVENTS] WebSocket closed: {status_code} {message}")
|
|
self.receive_event_queue.put_nowait(ValueError("WebSocket closed"))
|
|
|
|
ws = WebSocketApp(
|
|
ws_url,
|
|
on_message=on_message,
|
|
on_open=on_open,
|
|
on_close=on_close,
|
|
on_error=on_error,
|
|
)
|
|
|
|
from threading import Thread
|
|
|
|
wst = Thread(target=ws.run_forever)
|
|
wst.daemon = True
|
|
wst.start()
|
|
|
|
return ws
|
|
|
|
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)
|
|
|
|
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:
|
|
"""Check if an event has been seen recently."""
|
|
if event_id in self._seen_events:
|
|
return True
|
|
self._seen_events[event_id] = None
|
|
if len(self._seen_events) > MAX_SEEN_EVENTS:
|
|
self._seen_events.popitem(last=False)
|
|
return False
|
|
|
|
async def get_event(self):
|
|
"""Get the next relay message, consuming `OK` frames on the way.
|
|
|
|
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."""
|
|
self.subscription_id = "events-" + urlsafe_short_hash()[:32]
|
|
await self.send_req_queue.put(["REQ", self.subscription_id, *filters])
|
|
logger.info(
|
|
f"[EVENTS] Subscribed to NIP-52 events "
|
|
f"(sub: {self.subscription_id[:20]}...)"
|
|
)
|
|
|
|
async def unsubscribe(self):
|
|
"""Unsubscribe from current subscription."""
|
|
await self.send_req_queue.put(["CLOSE", self.subscription_id])
|
|
|
|
async def stop(self):
|
|
await self.unsubscribe()
|
|
self.running = False
|
|
await asyncio.sleep(2)
|
|
if self.ws:
|
|
try:
|
|
self.ws.close()
|
|
except Exception:
|
|
pass
|
|
self.ws = None
|