diff --git a/__init__.py b/__init__.py index a0b330d..6323038 100644 --- a/__init__.py +++ b/__init__.py @@ -34,6 +34,15 @@ scheduled_tasks: list[asyncio.Task] = [] # from nostr_hooks.publish_or_delete_nostr_event. nostr_client = None +# Reconciliation sweep for NIP-52 publishes that never reached a relay +# (aiolabs/events#35). Five minutes is well under the window in which a +# stale ticket count matters to a buyer, and the query costs nothing +# when there is no drift — the normal case returns zero rows. +REPUBLISH_SWEEP_INTERVAL = 300 +# Long enough for _start_nostr_client's own 10s wait plus the relay +# handshake, so the first pass isn't guaranteed to fail on a cold boot. +REPUBLISH_SWEEP_FIRST_DELAY = 60 + def events_stop(): for task in scheduled_tasks: @@ -117,5 +126,45 @@ def events_start(): task3 = create_permanent_unique_task("ext_events_nostr_sync", _sync_nostr_events) scheduled_tasks.append(task3) + async def _republish_pending_sweep(): + """Retry NIP-52 publishes that never landed. + + Inventory reaches clients only through the republished calendar + event, and a publish can fail (signer outage) or be skipped + entirely (no signer, no NostrClient) without anything noticing. + Both leave `nostr_publish_pending` set, so this sweep retries + from the DB rather than from an in-memory queue — it survives a + restart, which the previous behaviour did not. + + Quiet by design: on a healthy instance the query returns nothing + and this logs nothing. It only speaks up when there is drift. + """ + from .crud import get_events_pending_republish + from .nostr_hooks import publish_or_delete_nostr_event + + await asyncio.sleep(REPUBLISH_SWEEP_FIRST_DELAY) + while True: + try: + pending = await get_events_pending_republish() + if pending: + total = len(pending) + logger.info(f"[EVENTS] Republish sweep: {total} event(s) pending") + recovered = 0 + for event in pending: + take_down = event.canceled or event.status != "approved" + if await publish_or_delete_nostr_event(event, delete=take_down): + recovered += 1 + logger.info( + f"[EVENTS] Republish sweep: {recovered}/{total} recovered" + ) + except Exception as exc: + logger.error(f"[EVENTS] Republish sweep failed: {exc}") + await asyncio.sleep(REPUBLISH_SWEEP_INTERVAL) + + task4 = create_permanent_unique_task( + "ext_events_republish_sweep", _republish_pending_sweep + ) + scheduled_tasks.append(task4) + __all__ = ["db", "events_ext", "events_start", "events_static_files", "events_stop"] diff --git a/crud.py b/crud.py index bd5d3c1..caa0a2f 100644 --- a/crud.py +++ b/crud.py @@ -239,6 +239,26 @@ async def get_pending_events() -> list[Event]: ) +async def get_events_pending_republish() -> list[Event]: + """Events whose relay copy may be behind this row. + + `nostr_publish_pending` is set before every publish attempt and + cleared only on a confirmed success, so a row still flagged here + either failed to publish or never got the chance. Drives the + reconciliation sweep in `events_start`. + + Ordered oldest-first so a backlog drains in the order it accrued. + """ + return await db.fetchall( + """ + SELECT * FROM events.events + WHERE nostr_publish_pending = TRUE + ORDER BY time ASC + """, + model=Event, + ) + + async def get_settings() -> EventsSettings: """Singleton settings row, seeded by m010.""" row = await db.fetchone("SELECT * FROM events.settings WHERE id = 1") diff --git a/migrations_fork.py b/migrations_fork.py index ebe65a2..238a2d1 100644 --- a/migrations_fork.py +++ b/migrations_fork.py @@ -127,3 +127,32 @@ async def m002_ticket_payment_hash(db): "UPDATE events.ticket SET payment_hash = id " "WHERE payment_hash IS NULL OR payment_hash = ''" ) + + +async def m003_event_nostr_publish_pending(db): + """ + Add `events.nostr_publish_pending` — the marker that makes NIP-52 + publish drift queryable instead of invisible. + + Inventory reaches clients only through the republished calendar + event. When that publish doesn't land, the relay keeps serving the + counts it last saw and nothing anywhere records the divergence; it + has twice been caught only by a human reading a public page + (aiolabs/events#35, #51). + + The flag is set before each publish attempt and cleared only on a + confirmed success, so it covers *both* observed failure shapes: + an attempt that raised (a signer outage) and an attempt that was + never made at all (no signer resolved, no NostrClient). A periodic + sweep republishes whatever is still marked. + + Existing rows default to FALSE rather than TRUE: on upgrade we have + no evidence they're stale, and marking the whole table pending would + stampede the signer with a full-table republish on first boot. + `/republish-all` is the deliberate way to force that. + """ + await _alter_add_column_safe( + db, + "ALTER TABLE events.events " + "ADD COLUMN nostr_publish_pending BOOLEAN NOT NULL DEFAULT FALSE", + ) diff --git a/models.py b/models.py index b349794..8e7ec22 100644 --- a/models.py +++ b/models.py @@ -125,6 +125,10 @@ class Event(BaseModel): status: str = "approved" nostr_event_id: str | None = None nostr_event_created_at: int | None = None + # Set before every publish attempt, cleared on confirmed success. + # True means the relay's copy may be behind this row — see + # migrations_fork.m003 and the sweep in __init__.events_start. + nostr_publish_pending: bool = False @validator("categories", pre=True) def parse_categories(cls, v): diff --git a/nostr/nostr_client.py b/nostr/nostr_client.py index 4de332f..e0ae70f 100644 --- a/nostr/nostr_client.py +++ b/nostr/nostr_client.py @@ -19,6 +19,9 @@ 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 class NostrClient: @@ -78,17 +81,37 @@ class NostrClient: 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) - req = await self.send_req_queue.get() + 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: diff --git a/nostr_hooks.py b/nostr_hooks.py index 1d0d1c4..74cb8ec 100644 --- a/nostr_hooks.py +++ b/nostr_hooks.py @@ -12,7 +12,7 @@ from .models import Event from .nostr_publisher import publish_event_to_nostr -async def publish_or_delete_nostr_event(event: Event, *, delete: bool = False) -> None: +async def publish_or_delete_nostr_event(event: Event, *, delete: bool = False) -> bool: """Publish or delete the NIP-52 calendar event for `event`. Resolves a `NostrSigner` for the wallet owner — backend-agnostic @@ -22,7 +22,22 @@ async def publish_or_delete_nostr_event(event: Event, *, delete: bool = False) - `await signer.sign_event(...)` for signing. Failures are logged and swallowed so a Nostr outage doesn't break the HTTP flow that triggered the publish. + + Returns True when the event was signed and handed to the client, + False on any skip or failure. Callers are free to ignore it — the + `nostr_publish_pending` flag is the durable record, and the sweep + retries from that rather than from a return value. """ + # Mark before attempting, clear only on confirmed success. Doing it + # in this order is what makes "the attempt was never made" — no + # signer, no NostrClient, process died mid-flight — as visible as + # "the attempt raised". Cheap guard so a re-publish of an already + # pending row doesn't write twice; `set_ticket_paid` sets the flag + # inside its own update so the sale path adds no extra write. + if not event.nostr_publish_pending: + event.nostr_publish_pending = True + await update_event(event) + try: from lnbits.core.signers import resolve_for_wallet @@ -45,14 +60,22 @@ async def publish_or_delete_nostr_event(event: Event, *, delete: bool = False) - f"[EVENTS] No signer for wallet {event.wallet}, skipping " f"NIP-52 {'delete' if delete else 'publish'} for event {event.id}" ) - return + return False nostr_event = await publish_event_to_nostr( nostr_client, event, signer, delete=delete ) - if nostr_event and not delete: + if nostr_event is None: + return False + + event.nostr_publish_pending = False + if not delete: event.nostr_event_id = nostr_event.id event.nostr_event_created_at = nostr_event.created_at - await update_event(event) + await update_event(event) + return True except Exception as exc: - logger.warning(f"[EVENTS] Nostr publish failed: {exc}") + # ERROR, not warning: the row stays flagged and its published + # counts stay behind until the sweep or a later edit succeeds. + logger.error(f"[EVENTS] Nostr publish failed for event {event.id}: {exc}") + return False diff --git a/nostr_publisher.py b/nostr_publisher.py index c961b80..485acd4 100644 --- a/nostr_publisher.py +++ b/nostr_publisher.py @@ -212,5 +212,8 @@ async def publish_event_to_nostr( return nostr_event except Exception as e: - logger.warning(f"[EVENTS] Failed to publish to Nostr: {e}") + # ERROR, not warning: this is the signer-outage shape of + # aiolabs/events#35 — the calendar event never reaches the relay + # and the published ticket counts stop tracking the DB. + logger.error(f"[EVENTS] Failed to publish event {event.id} to Nostr: {e}") return None diff --git a/services.py b/services.py index f69b76e..311b7e2 100644 --- a/services.py +++ b/services.py @@ -71,6 +71,10 @@ async def set_ticket_paid(ticket: Ticket) -> Ticket: assert event, "Couldn't get event from ticket being paid" event.sold += 1 event.amount_tickets -= 1 + # Flag inside this same write: the counters and "the relay does + # not know about them yet" land atomically, so a crash between + # here and the publish still leaves the drift discoverable. + event.nostr_publish_pending = True await update_event(event) # Republish the NIP-52 calendar event so connected clients see diff --git a/tests/test_nostr_publish_pending.py b/tests/test_nostr_publish_pending.py new file mode 100644 index 0000000..3caace0 --- /dev/null +++ b/tests/test_nostr_publish_pending.py @@ -0,0 +1,171 @@ +"""The `nostr_publish_pending` marker and its lifecycle. + +Inventory reaches clients only through the republished NIP-52 calendar +event. These tests pin the invariant that makes drift recoverable: the +flag goes up before every attempt and comes down only on a confirmed +success, so every shape of failure — raised, skipped, never attempted — +leaves the row queryable by the sweep. +""" + +from datetime import datetime, timezone +from types import SimpleNamespace +from unittest.mock import AsyncMock + +import pytest + +from .. import nostr_hooks, services +from ..models import Event, Ticket + + +def _event(**kwargs) -> Event: + defaults = { + "id": "evt", + "wallet": "w", + "name": "Test", + "info": "", + "closing_date": "2030-01-01", + "event_start_date": "2030-01-01", + "event_end_date": "2030-01-02", + "currency": "sat", + "price_per_ticket": 1000, + "amount_tickets": 10, + "time": datetime.now(timezone.utc), + "status": "approved", + } + defaults.update(kwargs) + return Event(**defaults) + + +@pytest.fixture +def saved(monkeypatch): + """Capture every update_event write so ordering can be asserted.""" + writes: list[bool] = [] + + async def _update(event): + writes.append(event.nostr_publish_pending) + return event + + monkeypatch.setattr(nostr_hooks, "update_event", _update) + return writes + + +def _signer(monkeypatch, signer): + monkeypatch.setattr( + "lnbits.core.signers.resolve_for_wallet", AsyncMock(return_value=signer) + ) + + +def _publisher(monkeypatch, result): + monkeypatch.setattr( + nostr_hooks, "publish_event_to_nostr", AsyncMock(return_value=result) + ) + + +@pytest.mark.asyncio +async def test_success_raises_then_clears_the_flag(monkeypatch, saved): + event = _event() + _signer(monkeypatch, SimpleNamespace(pubkey="pk")) + _publisher(monkeypatch, SimpleNamespace(id="nid", created_at=123)) + + assert await nostr_hooks.publish_or_delete_nostr_event(event) is True + # Flagged before the attempt, cleared after it — in that order. + assert saved == [True, False] + assert event.nostr_publish_pending is False + assert event.nostr_event_id == "nid" + assert event.nostr_event_created_at == 123 + + +@pytest.mark.asyncio +async def test_missing_signer_leaves_the_flag_up(monkeypatch, saved): + event = _event() + _signer(monkeypatch, None) + + assert await nostr_hooks.publish_or_delete_nostr_event(event) is False + assert saved == [True] + assert event.nostr_publish_pending is True + + +@pytest.mark.asyncio +async def test_publisher_returning_none_leaves_the_flag_up(monkeypatch, saved): + """The no-NostrClient shape: nothing raised, nothing published.""" + event = _event() + _signer(monkeypatch, SimpleNamespace(pubkey="pk")) + _publisher(monkeypatch, None) + + assert await nostr_hooks.publish_or_delete_nostr_event(event) is False + assert saved == [True] + assert event.nostr_publish_pending is True + + +@pytest.mark.asyncio +async def test_raised_publish_leaves_the_flag_up(monkeypatch, saved): + event = _event() + _signer(monkeypatch, SimpleNamespace(pubkey="pk")) + monkeypatch.setattr( + nostr_hooks, + "publish_event_to_nostr", + AsyncMock(side_effect=RuntimeError("signer timeout")), + ) + + assert await nostr_hooks.publish_or_delete_nostr_event(event) is False + assert saved == [True] + assert event.nostr_publish_pending is True + + +@pytest.mark.asyncio +async def test_already_pending_row_is_not_re_flagged(monkeypatch, saved): + """The sweep re-publishing a flagged row writes once, not twice.""" + event = _event(nostr_publish_pending=True) + _signer(monkeypatch, SimpleNamespace(pubkey="pk")) + _publisher(monkeypatch, SimpleNamespace(id="nid", created_at=123)) + + assert await nostr_hooks.publish_or_delete_nostr_event(event) is True + assert saved == [False] + + +@pytest.mark.asyncio +async def test_delete_clears_the_flag_without_touching_the_coordinate( + monkeypatch, saved +): + """A take-down must not overwrite the id/created_at of the event it + just deleted — the kind-5 has its own.""" + event = _event(nostr_event_id="old", nostr_event_created_at=100) + _signer(monkeypatch, SimpleNamespace(pubkey="pk")) + _publisher(monkeypatch, SimpleNamespace(id="del", created_at=999)) + + assert await nostr_hooks.publish_or_delete_nostr_event(event, delete=True) is True + assert event.nostr_publish_pending is False + assert event.nostr_event_id == "old" + assert event.nostr_event_created_at == 100 + + +@pytest.mark.asyncio +async def test_sale_flags_the_event_in_the_same_write(monkeypatch): + """`set_ticket_paid` must flag inside its own update, so the counters + and "the relay doesn't know yet" land atomically.""" + event = _event(sold=4, amount_tickets=6) + seen: list[tuple[int, int, bool]] = [] + + async def _update_event(ev): + seen.append((ev.sold, ev.amount_tickets, ev.nostr_publish_pending)) + return ev + + monkeypatch.setattr(services, "update_ticket", AsyncMock()) + monkeypatch.setattr(services, "get_event", AsyncMock(return_value=event)) + monkeypatch.setattr(services, "update_event", _update_event) + monkeypatch.setattr(services, "publish_or_delete_nostr_event", AsyncMock()) + + ticket = Ticket( + id="t1", + wallet="w", + event="evt", + name="A", + email="a@example.com", + registered=False, + paid=False, + time=datetime.now(timezone.utc), + reg_timestamp=datetime.now(timezone.utc), + ) + await services.set_ticket_paid(ticket) + + assert seen == [(5, 5, True)]