A flagged row recovers on its own instead of waiting for the next sale that happens to land while the signer is healthy — or for an operator who already knows to run /republish-all, which was the only recovery path and requires knowing about drift that nothing reported. Retrying from the DB rather than an in-memory queue means the retry survives a restart, and it needs no theory about why the publish didn't land: the sweep covers the signer outage of #35 and the silent skip of #51 identically, along with causes nobody has hit yet. Runs every 5 minutes, take-down branch mirroring the publish/delete split the CRUD endpoints already use. Quiet by design — on a healthy instance the query returns nothing and it logs nothing. Refs #35
170 lines
6 KiB
Python
170 lines
6 KiB
Python
import asyncio
|
|
|
|
from fastapi import APIRouter
|
|
from loguru import logger
|
|
|
|
from .crud import db
|
|
from .tasks import wait_for_paid_invoices
|
|
from .views import events_generic_router
|
|
from .views_api import (
|
|
events_api_router,
|
|
promo_api_router,
|
|
qr_api_router,
|
|
tickets_api_router,
|
|
)
|
|
|
|
events_ext: APIRouter = APIRouter(prefix="/events", tags=["Events"])
|
|
events_ext.include_router(events_generic_router)
|
|
events_ext.include_router(events_api_router)
|
|
events_ext.include_router(tickets_api_router)
|
|
events_ext.include_router(qr_api_router)
|
|
events_ext.include_router(promo_api_router)
|
|
|
|
events_static_files = [
|
|
{
|
|
"path": "/events/static",
|
|
"name": "events_static",
|
|
}
|
|
]
|
|
|
|
scheduled_tasks: list[asyncio.Task] = []
|
|
|
|
# Module-level NostrClient — None when nostrclient is unavailable. Set by the
|
|
# bootstrap task in events_start() and read via dynamic attribute lookup
|
|
# 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:
|
|
try:
|
|
task.cancel()
|
|
except Exception as ex:
|
|
logger.warning(ex)
|
|
|
|
global nostr_client
|
|
if nostr_client:
|
|
asyncio.get_event_loop().create_task(nostr_client.stop())
|
|
|
|
|
|
def events_start():
|
|
from lnbits.tasks import create_permanent_unique_task
|
|
|
|
task1 = create_permanent_unique_task("ext_events", wait_for_paid_invoices)
|
|
scheduled_tasks.append(task1)
|
|
|
|
# Register nostr-transport RPCs. Swallow ImportError on older LNbits
|
|
# versions that pre-date the transport (the events extension still
|
|
# works fine via HTTP without it).
|
|
try:
|
|
from lnbits.core.services.nostr_transport.dispatcher import (
|
|
AUTH_WALLET,
|
|
register_rpc,
|
|
)
|
|
|
|
from .transport_rpcs import (
|
|
handle_events_list_event_tickets,
|
|
handle_events_ticket_register,
|
|
)
|
|
|
|
register_rpc(
|
|
"events_ticket_register", handle_events_ticket_register, AUTH_WALLET
|
|
)
|
|
register_rpc(
|
|
"events_list_event_tickets",
|
|
handle_events_list_event_tickets,
|
|
AUTH_WALLET,
|
|
)
|
|
logger.info(
|
|
"[EVENTS] Registered nostr-transport RPCs: "
|
|
"events_ticket_register, events_list_event_tickets"
|
|
)
|
|
except ImportError:
|
|
logger.info(
|
|
"[EVENTS] nostr_transport not available on this LNbits — "
|
|
"ticket scanner over Nostr disabled, HTTP endpoint still works"
|
|
)
|
|
|
|
async def _start_nostr_client():
|
|
global nostr_client
|
|
await asyncio.sleep(10) # Wait for nostrclient to be ready
|
|
try:
|
|
from .nostr.nostr_client import NostrClient
|
|
|
|
nostr_client = NostrClient()
|
|
logger.info("[EVENTS] Starting NostrClient for NIP-52 sync")
|
|
await nostr_client.run_forever()
|
|
except Exception as exc:
|
|
logger.warning(f"[EVENTS] NostrClient failed to start: {exc}")
|
|
logger.info("[EVENTS] Events will work without Nostr sync")
|
|
|
|
task2 = create_permanent_unique_task("ext_events_nostr", _start_nostr_client)
|
|
scheduled_tasks.append(task2)
|
|
|
|
async def _sync_nostr_events():
|
|
global nostr_client
|
|
await asyncio.sleep(15) # Wait for NostrClient to connect
|
|
if not nostr_client:
|
|
logger.info("[EVENTS] No NostrClient, skipping Nostr sync")
|
|
return
|
|
try:
|
|
from .nostr_sync import wait_for_nostr_events
|
|
|
|
await wait_for_nostr_events(nostr_client)
|
|
except Exception as exc:
|
|
logger.error(f"[EVENTS] Nostr sync task failed: {exc}")
|
|
|
|
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"]
|