feat(nostr): make publish drift queryable and self-healing #55
2 changed files with 69 additions and 0 deletions
feat(nostr): sweep republishes events flagged as pending
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
commit
643322131e
49
__init__.py
49
__init__.py
|
|
@ -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
20
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:
|
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")
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue