feat(nostr): make publish drift queryable and self-healing #55

Merged
padreug merged 3 commits from feat/nostr-publish-reconciliation into main 2026-09-26 21:47:41 +00:00
9 changed files with 333 additions and 7 deletions

View file

@ -34,6 +34,15 @@ scheduled_tasks: list[asyncio.Task] = []
# from nostr_hooks.publish_or_delete_nostr_event. # from nostr_hooks.publish_or_delete_nostr_event.
nostr_client = None 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(): def events_stop():
for task in scheduled_tasks: for task in scheduled_tasks:
@ -117,5 +126,45 @@ def events_start():
task3 = create_permanent_unique_task("ext_events_nostr_sync", _sync_nostr_events) task3 = create_permanent_unique_task("ext_events_nostr_sync", _sync_nostr_events)
scheduled_tasks.append(task3) 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"] __all__ = ["db", "events_ext", "events_start", "events_static_files", "events_stop"]

20
crud.py
View file

@ -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: async def get_settings() -> EventsSettings:
"""Singleton settings row, seeded by m010.""" """Singleton settings row, seeded by m010."""
row = await db.fetchone("SELECT * FROM events.settings WHERE id = 1") row = await db.fetchone("SELECT * FROM events.settings WHERE id = 1")

View file

@ -127,3 +127,32 @@ async def m002_ticket_payment_hash(db):
"UPDATE events.ticket SET payment_hash = id " "UPDATE events.ticket SET payment_hash = id "
"WHERE payment_hash IS NULL OR payment_hash = ''" "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",
)

View file

@ -125,6 +125,10 @@ class Event(BaseModel):
status: str = "approved" status: str = "approved"
nostr_event_id: str | None = None nostr_event_id: str | None = None
nostr_event_created_at: int | 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) @validator("categories", pre=True)
def parse_categories(cls, v): def parse_categories(cls, v):

View file

@ -19,6 +19,9 @@ from websocket import WebSocketApp
from .event import NostrEvent from .event import NostrEvent
MAX_SEEN_EVENTS = 500 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: class NostrClient:
@ -78,17 +81,37 @@ class NostrClient:
async def run_forever(self): async def run_forever(self):
self.running = True 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: while self.running:
try: try:
if not self.is_websocket_connected: if not self.is_websocket_connected:
self.ws = await self.connect() self.ws = await self.connect()
await asyncio.sleep(5) await asyncio.sleep(5)
if held_req is not None:
req = held_req
else:
req = await self.send_req_queue.get() req = await self.send_req_queue.get()
held_req, held_attempts = req, held_attempts + 1
assert self.ws assert self.ws
self.ws.send(json.dumps(req)) self.ws.send(json.dumps(req))
held_req, held_attempts = None, 0
except Exception as ex: except Exception as ex:
logger.warning(f"[EVENTS] NostrClient error: {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) await asyncio.sleep(60)
def is_duplicate_event(self, event_id: str) -> bool: def is_duplicate_event(self, event_id: str) -> bool:

View file

@ -12,7 +12,7 @@ from .models import Event
from .nostr_publisher import publish_event_to_nostr 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`. """Publish or delete the NIP-52 calendar event for `event`.
Resolves a `NostrSigner` for the wallet owner — backend-agnostic 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 `await signer.sign_event(...)` for signing. Failures are logged
and swallowed so a Nostr outage doesn't break the HTTP flow that and swallowed so a Nostr outage doesn't break the HTTP flow that
triggered the publish. 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: try:
from lnbits.core.signers import resolve_for_wallet 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"[EVENTS] No signer for wallet {event.wallet}, skipping "
f"NIP-52 {'delete' if delete else 'publish'} for event {event.id}" f"NIP-52 {'delete' if delete else 'publish'} for event {event.id}"
) )
return return False
nostr_event = await publish_event_to_nostr( nostr_event = await publish_event_to_nostr(
nostr_client, event, signer, delete=delete 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_id = nostr_event.id
event.nostr_event_created_at = nostr_event.created_at event.nostr_event_created_at = nostr_event.created_at
await update_event(event) await update_event(event)
return True
except Exception as exc: 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

View file

@ -212,5 +212,8 @@ async def publish_event_to_nostr(
return nostr_event return nostr_event
except Exception as e: 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 return None

View file

@ -71,6 +71,10 @@ async def set_ticket_paid(ticket: Ticket) -> Ticket:
assert event, "Couldn't get event from ticket being paid" assert event, "Couldn't get event from ticket being paid"
event.sold += 1 event.sold += 1
event.amount_tickets -= 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) await update_event(event)
# Republish the NIP-52 calendar event so connected clients see # Republish the NIP-52 calendar event so connected clients see

View file

@ -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)]