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
89 lines
2.9 KiB
Python
89 lines
2.9 KiB
Python
import asyncio
|
|
import threading
|
|
|
|
from loguru import logger
|
|
|
|
from .crud import get_relays
|
|
from .nostr.message_pool import (
|
|
CommandResultMessage,
|
|
EndOfStoredEventsMessage,
|
|
EventMessage,
|
|
NoticeMessage,
|
|
)
|
|
from .router import NostrRouter, all_routers, nostr_client
|
|
|
|
|
|
async def init_relays():
|
|
# get relays from db
|
|
relays = await get_relays()
|
|
# set relays and connect to them
|
|
valid_relays = [r.url for r in relays if r.url]
|
|
|
|
nostr_client.reconnect(valid_relays)
|
|
|
|
|
|
async def check_relays():
|
|
"""Check relays that have been disconnected"""
|
|
while True:
|
|
try:
|
|
await asyncio.sleep(20)
|
|
nostr_client.relay_manager.check_and_restart_relays()
|
|
except Exception as e:
|
|
logger.warning(f"Cannot restart relays: '{e!s}'.")
|
|
|
|
|
|
async def subscribe_events():
|
|
while not [r.connected for r in nostr_client.relay_manager.relays.values()]:
|
|
await asyncio.sleep(2)
|
|
|
|
def callback_events(event_message: EventMessage):
|
|
sub_id = event_message.subscription_id
|
|
if sub_id not in NostrRouter.received_subscription_events:
|
|
NostrRouter.received_subscription_events[sub_id] = [event_message]
|
|
return
|
|
|
|
# do not add duplicate events (by event id)
|
|
ids = [e.event_id for e in NostrRouter.received_subscription_events[sub_id]]
|
|
if event_message.event_id in ids:
|
|
return
|
|
|
|
NostrRouter.received_subscription_events[sub_id].append(event_message)
|
|
|
|
def callback_notices(notice_message: NoticeMessage):
|
|
if notice_message not in NostrRouter.received_subscription_notices:
|
|
NostrRouter.received_subscription_notices.append(notice_message)
|
|
|
|
def callback_eose_notices(event_message: EndOfStoredEventsMessage):
|
|
sub_id = event_message.subscription_id
|
|
if sub_id in NostrRouter.received_subscription_eosenotices:
|
|
return
|
|
|
|
NostrRouter.received_subscription_eosenotices[sub_id] = event_message
|
|
|
|
def callback_command_results(result_message: CommandResultMessage):
|
|
event_id = result_message.event_id
|
|
# Only keep OKs some client is still waiting on; the rest would
|
|
# accumulate forever (events published by other relay users, or
|
|
# results arriving after the client already got its reply).
|
|
if not any(event_id in r.pending_publishes for r in all_routers):
|
|
return
|
|
NostrRouter.received_command_results.setdefault(event_id, []).append(
|
|
result_message
|
|
)
|
|
|
|
def wrap_async_subscribe():
|
|
asyncio.run(
|
|
nostr_client.subscribe(
|
|
callback_events,
|
|
callback_notices,
|
|
callback_eose_notices,
|
|
callback_command_results,
|
|
)
|
|
)
|
|
|
|
t = threading.Thread(
|
|
target=wrap_async_subscribe,
|
|
name="Nostr-event-subscription",
|
|
daemon=True,
|
|
)
|
|
t.start()
|