Compare commits

..

No commits in common. "13eb54906640aeedd5d75e028b2a3d993da229ab" and "cac7d16ece5f354248531aab1524f7d4059bca7d" have entirely different histories.

3 changed files with 9 additions and 200 deletions

View file

@ -22,11 +22,6 @@ MAX_SEEN_EVENTS = 500
# How many times a dequeued req is retried before it is dropped. Bounded # 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. # so one unsendable message can't block every later publish behind it.
MAX_SEND_ATTEMPTS = 3 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: class NostrClient:
@ -37,10 +32,6 @@ class NostrClient:
self.subscription_id = "events-" + urlsafe_short_hash()[:32] self.subscription_id = "events-" + urlsafe_short_hash()[:32]
self.running = False self.running = False
self._seen_events: OrderedDict[str, None] = OrderedDict() 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 @property
def is_websocket_connected(self): def is_websocket_connected(self):
@ -133,82 +124,14 @@ class NostrClient:
return False return False
async def get_event(self): async def get_event(self):
"""Get the next relay message, consuming `OK` frames on the way. """Get next event from the receive queue."""
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() value = await self.receive_event_queue.get()
if isinstance(value, ValueError): if isinstance(value, ValueError):
self._fail_pending_oks("connection closed")
raise value raise value
if self._settle_ok(value):
continue
return value return value
def _settle_ok(self, message) -> bool: async def publish_nostr_event(self, e: NostrEvent):
"""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()]) 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]): async def subscribe(self, filters: list[dict]):
"""Subscribe to events matching the given filters.""" """Subscribe to events matching the given filters."""

View file

@ -203,16 +203,7 @@ async def publish_event_to_nostr(
nostr_event.pubkey = signed["pubkey"] nostr_event.pubkey = signed["pubkey"]
nostr_event.sig = signed["sig"] nostr_event.sig = signed["sig"]
accepted = await nostr_client.publish_nostr_event(nostr_event) await nostr_client.publish_nostr_event(nostr_event)
if not accepted:
# Returning None keeps `nostr_publish_pending` set, so the
# sweep retries instead of recording a delivery that never
# happened (aiolabs/events#56).
logger.warning(
f"[EVENTS] Relay did not confirm NIP-52 "
f"{'delete' if delete else 'calendar'} event for {event.id}"
)
return None
logger.info( logger.info(
f"[EVENTS] Published NIP-52 {'delete' if delete else 'calendar'} " f"[EVENTS] Published NIP-52 {'delete' if delete else 'calendar'} "

View file

@ -1,105 +0,0 @@
"""Publish confirmation against the relay's `OK` (aiolabs/events#56).
Queueing is not delivery. nostrclient drops an EVENT outright when no
relay is connected and answers `OK false`; before this, that reply was
discarded and the publish reported success, which once cleared the
`nostr_publish_pending` flag on a republish that never left the
building.
"""
import asyncio
import json
import pytest
from ..nostr import nostr_client as nc
from ..nostr.event import NostrEvent
def _event(event_id: str = "a" * 64) -> NostrEvent:
e = NostrEvent(pubkey="b" * 64, created_at=0, kind=31923)
e.id = event_id
return e
async def _publish_and_reply(client, event, reply, delay=0.01):
"""Start a publish, then feed `reply` through the receive path."""
task = asyncio.create_task(client.publish_nostr_event(event))
await asyncio.sleep(delay) # let the future register
if reply is not None:
client.receive_event_queue.put_nowait(reply)
consumer = asyncio.create_task(client.get_event())
await asyncio.sleep(delay)
consumer.cancel()
return await task
@pytest.mark.asyncio
async def test_accepted_publish_returns_true():
client = nc.NostrClient()
event = _event()
ok = json.dumps(["OK", event.id, True, ""])
assert await _publish_and_reply(client, event, ok) is True
assert client._pending_oks == {}
@pytest.mark.asyncio
async def test_rejected_publish_returns_false():
"""The shape that bit us: no relay connected, so nostrclient's
router answers `OK false` without the event ever being sent."""
client = nc.NostrClient()
event = _event()
ok = json.dumps(["OK", event.id, False, "error: no relays connected"])
assert await _publish_and_reply(client, event, ok) is False
assert client._pending_oks == {}
@pytest.mark.asyncio
async def test_missing_ok_times_out_as_unconfirmed(monkeypatch):
monkeypatch.setattr(nc, "PUBLISH_OK_TIMEOUT_SECONDS", 0.05)
client = nc.NostrClient()
assert await _publish_and_reply(client, _event(), None) is False
assert client._pending_oks == {}
@pytest.mark.asyncio
async def test_ok_for_a_different_event_does_not_settle_ours(monkeypatch):
monkeypatch.setattr(nc, "PUBLISH_OK_TIMEOUT_SECONDS", 0.05)
client = nc.NostrClient()
other = json.dumps(["OK", "c" * 64, True, ""])
assert await _publish_and_reply(client, _event(), other) is False
@pytest.mark.asyncio
async def test_get_event_swallows_ok_and_forwards_everything_else():
client = nc.NostrClient()
client.receive_event_queue.put_nowait(json.dumps(["OK", "d" * 64, True, ""]))
forwarded = json.dumps(["EVENT", "sub", {"id": "e" * 64}])
client.receive_event_queue.put_nowait(forwarded)
assert await client.get_event() == forwarded
@pytest.mark.asyncio
async def test_disconnect_settles_inflight_publishes_immediately(monkeypatch):
"""A dropped socket must not leave the caller waiting the full
timeout for an OK that can no longer arrive."""
monkeypatch.setattr(nc, "PUBLISH_OK_TIMEOUT_SECONDS", 30)
client = nc.NostrClient()
event = _event()
task = asyncio.create_task(client.publish_nostr_event(event))
await asyncio.sleep(0.01)
client.receive_event_queue.put_nowait(ValueError("WebSocket closed"))
consumer = asyncio.create_task(client.get_event())
await asyncio.sleep(0.01)
consumer.cancel()
assert await asyncio.wait_for(task, 1) is False
def test_settle_ok_ignores_non_ok_frames():
client = nc.NostrClient()
assert client._settle_ok(json.dumps(["EVENT", "sub", {}])) is False
assert client._settle_ok(json.dumps(["EOSE", "sub"])) is False
assert client._settle_ok("not json") is False
assert client._settle_ok(json.dumps(["OK", "f" * 64, True, ""])) is True