nostrclient never answered a client's EVENT with the NIP-01 `["OK", <id>, <accepted>, <message>]` command result. Clients built on nostr-tools and similar libraries wait for that reply before treating a publish as successful, so NWC wallet apps paired through the public endpoint reported "publish failed" even though the request had been fanned out and answered. Relay OKs now flow through the message pool like events and notices. The router tracks each EVENT a client publishes and replies exactly once: `true` as soon as any relay accepts, `false` once every relay connected at publish time has rejected it, after a 10 s timeout, or immediately when no relay is connected. OKs nobody is waiting on are dropped at the pump so the shared result map cannot grow unbounded. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_013Tbyw6FwjhEJg3gHfPHxWt
86 lines
2.8 KiB
Python
86 lines
2.8 KiB
Python
import asyncio
|
|
|
|
from loguru import logger
|
|
|
|
from ..relay_manager import RelayManager
|
|
|
|
|
|
class NostrClient:
|
|
relay_manager: RelayManager
|
|
running: bool
|
|
|
|
def __init__(self):
|
|
self.running = True
|
|
self.relay_manager = RelayManager()
|
|
|
|
def connect(self, relays):
|
|
for relay in relays:
|
|
try:
|
|
self.relay_manager.add_relay(relay)
|
|
except Exception as e:
|
|
logger.debug(e)
|
|
self.running = True
|
|
|
|
def reconnect(self, relays):
|
|
self.relay_manager.remove_relays()
|
|
self.connect(relays)
|
|
|
|
def close(self):
|
|
try:
|
|
self.relay_manager.close_all_subscriptions()
|
|
self.relay_manager.close_connections()
|
|
|
|
self.running = False
|
|
except Exception as e:
|
|
logger.error(e)
|
|
|
|
async def subscribe(
|
|
self,
|
|
callback_events_func=None,
|
|
callback_notices_func=None,
|
|
callback_eosenotices_func=None,
|
|
callback_command_results_func=None,
|
|
):
|
|
while self.running:
|
|
self._check_events(callback_events_func)
|
|
self._check_notices(callback_notices_func)
|
|
self._check_eos_notices(callback_eosenotices_func)
|
|
self._check_command_results(callback_command_results_func)
|
|
|
|
await asyncio.sleep(0.2)
|
|
|
|
def _check_events(self, callback_events_func=None):
|
|
try:
|
|
while self.relay_manager.message_pool.has_events():
|
|
event_msg = self.relay_manager.message_pool.get_event()
|
|
if callback_events_func:
|
|
callback_events_func(event_msg)
|
|
except Exception as e:
|
|
logger.debug(e)
|
|
|
|
def _check_notices(self, callback_notices_func=None):
|
|
try:
|
|
while self.relay_manager.message_pool.has_notices():
|
|
event_msg = self.relay_manager.message_pool.get_notice()
|
|
if callback_notices_func:
|
|
callback_notices_func(event_msg)
|
|
except Exception as e:
|
|
logger.debug(e)
|
|
|
|
def _check_eos_notices(self, callback_eosenotices_func=None):
|
|
try:
|
|
while self.relay_manager.message_pool.has_eose_notices():
|
|
event_msg = self.relay_manager.message_pool.get_eose_notice()
|
|
if callback_eosenotices_func:
|
|
callback_eosenotices_func(event_msg)
|
|
except Exception as e:
|
|
logger.debug(e)
|
|
|
|
def _check_command_results(self, callback_command_results_func=None):
|
|
try:
|
|
while self.relay_manager.message_pool.has_command_results():
|
|
result_msg = self.relay_manager.message_pool.get_command_result()
|
|
if callback_command_results_func:
|
|
callback_command_results_func(result_msg)
|
|
except Exception as e:
|
|
logger.debug(e)
|