nostrclient/tasks.py
Padreug 7aa0ce14db feat: reply OK to client-published events
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
2026-09-12 13:31:00 +02:00

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()