Compare commits
No commits in common. "114f3406a6f2efefc8692aed84f8737e46a0fe8e" and "d2b8550d7fe096f8817f79a93119cc9e3df279cd" have entirely different histories.
114f3406a6
...
d2b8550d7f
9 changed files with 7 additions and 333 deletions
49
__init__.py
49
__init__.py
|
|
@ -34,15 +34,6 @@ 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:
|
||||
|
|
@ -126,45 +117,5 @@ 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"]
|
||||
|
|
|
|||
20
crud.py
20
crud.py
|
|
@ -239,26 +239,6 @@ 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")
|
||||
|
|
|
|||
|
|
@ -127,32 +127,3 @@ 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",
|
||||
)
|
||||
|
|
|
|||
|
|
@ -125,10 +125,6 @@ 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):
|
||||
|
|
|
|||
|
|
@ -19,9 +19,6 @@ 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:
|
||||
|
|
@ -81,37 +78,17 @@ 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)
|
||||
|
||||
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
|
||||
req = await self.send_req_queue.get()
|
||||
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:
|
||||
|
|
|
|||
|
|
@ -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) -> bool:
|
||||
async def publish_or_delete_nostr_event(event: Event, *, delete: bool = False) -> None:
|
||||
"""Publish or delete the NIP-52 calendar event for `event`.
|
||||
|
||||
Resolves a `NostrSigner` for the wallet owner — backend-agnostic
|
||||
|
|
@ -22,22 +22,7 @@ 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
|
||||
|
||||
|
|
@ -60,22 +45,14 @@ 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 False
|
||||
return
|
||||
|
||||
nostr_event = await publish_event_to_nostr(
|
||||
nostr_client, event, signer, delete=delete
|
||||
)
|
||||
if nostr_event is None:
|
||||
return False
|
||||
|
||||
event.nostr_publish_pending = False
|
||||
if not delete:
|
||||
if nostr_event and not delete:
|
||||
event.nostr_event_id = nostr_event.id
|
||||
event.nostr_event_created_at = nostr_event.created_at
|
||||
await update_event(event)
|
||||
return True
|
||||
await update_event(event)
|
||||
except Exception as 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
|
||||
logger.warning(f"[EVENTS] Nostr publish failed: {exc}")
|
||||
|
|
|
|||
|
|
@ -212,8 +212,5 @@ async def publish_event_to_nostr(
|
|||
return nostr_event
|
||||
|
||||
except Exception as 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}")
|
||||
logger.warning(f"[EVENTS] Failed to publish to Nostr: {e}")
|
||||
return None
|
||||
|
|
|
|||
|
|
@ -71,10 +71,6 @@ 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
|
||||
|
|
|
|||
|
|
@ -1,171 +0,0 @@
|
|||
"""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)]
|
||||
Loading…
Add table
Add a link
Reference in a new issue