Some checks failed
lint.yml / feat(nostr): publish the active ticket wave, not the roll-up (pull_request) Failing after 0s
Since upstream v1.6.8 price, currency and inventory belong to time-boxed ticket waves, and the event-level fields `sync_event_ticket_waves` derives are the PRIMARY wave's price/currency and the SUM of every wave's stock. The NIP-52 publisher read those, so as soon as an organiser created a second wave the public card would advertise the early-bird price after early bird closed and count stock in waves that had not opened. Refs #61. `build_nip52_event` now describes the wave a buyer can actually buy from: - several waves can be open at once, and a publisher has no one to ask which one the buyer wants (the purchase endpoint errors with "Please select a ticket wave"), so it advertises the CHEAPEST open wave — the price a buyer is able to obtain. Deviation recorded in docs/upstream-candidates.md. - with no open wave, `tickets_available` is 0 and never omitted: omission used to mean "unlimited", which #34/#62 removed as a concept. - `tickets_payment_methods` is scoped to the advertised wave too. It was derived from `event.allow_fiat` — the primary wave's — so it could offer a fiat rail while `tickets_allow_fiat` was absent and the purchase endpoint would refuse it. They are the same fact and now come from the same place. Wave boundaries are time-driven, and every republish we have is sale-driven, so nothing fires when early bird ends at midnight. Rather than add a scheduler, a publish records which wave it advertised (`nostr_published_wave_id`, m004) and the reconciliation sweep compares that against the wave that would be advertised now, setting `nostr_publish_pending` on a mismatch — reusing the existing retry path. NULL means "never published", which the sweep leaves alone so an upgrade does not republish the whole table on first boot. The selection rule lives in `models.advertised_ticket_wave` so the publisher and the drift detector cannot disagree about what is on the relay. 17 new tests; 114 pass. ruff, black, prettier clean; mypy error set still identical to HEAD's baseline.
396 lines
13 KiB
Python
396 lines
13 KiB
Python
import json
|
|
from datetime import datetime, timedelta, timezone
|
|
from typing import cast
|
|
|
|
from lnbits.db import Database, Filters, Page
|
|
from lnbits.helpers import urlsafe_short_hash
|
|
|
|
from .models import (
|
|
CreateEvent,
|
|
Event,
|
|
EventsSettings,
|
|
Ticket,
|
|
TicketExtra,
|
|
TicketFilters,
|
|
advertised_wave_key,
|
|
sync_event_ticket_waves,
|
|
)
|
|
|
|
db = Database("ext_events")
|
|
|
|
|
|
def _parse_ticket_row(row) -> dict:
|
|
"""Normalize a ticket row before constructing a Ticket model.
|
|
|
|
- Empty-string sentinels in name/email (used because the DB columns are
|
|
NOT NULL but the Pydantic field is Optional when user_id is set) are
|
|
converted back to None.
|
|
- The `extra` JSON column may come back as a string when the row is
|
|
fetched without a model= argument; parse it so Pydantic can build
|
|
TicketExtra from a dict.
|
|
"""
|
|
ticket_data = dict(row)
|
|
|
|
if ticket_data.get("name") == "":
|
|
ticket_data["name"] = None
|
|
if ticket_data.get("email") == "":
|
|
ticket_data["email"] = None
|
|
|
|
extra = ticket_data.get("extra")
|
|
if isinstance(extra, str):
|
|
ticket_data["extra"] = json.loads(extra) if extra else {}
|
|
|
|
return ticket_data
|
|
|
|
|
|
async def create_ticket(
|
|
payment_hash: str,
|
|
wallet: str,
|
|
event: str,
|
|
name: str | None = None,
|
|
email: str | None = None,
|
|
user_id: str | None = None,
|
|
extra: dict | None = None,
|
|
ticket_id: str | None = None,
|
|
) -> Ticket:
|
|
"""Persist one ticket row.
|
|
|
|
`payment_hash` is the LNbits invoice hash shared across all rows
|
|
of a multi-ticket purchase. `ticket_id` is the row primary key /
|
|
scannable id; defaults to `payment_hash` for single-ticket
|
|
purchases so the legacy id == payment_hash invariant holds.
|
|
Multi-ticket callers pass a unique uuid here so each attendee
|
|
gets a distinct scannable QR.
|
|
"""
|
|
now = datetime.now(timezone.utc)
|
|
row_id = ticket_id or payment_hash
|
|
|
|
# name/email columns are NOT NULL in the schema, so we store "" when a
|
|
# value is absent. _parse_ticket_row reverses this on read. A user_id
|
|
# ticket may carry an email too — that is how logged-in webapp buyers get
|
|
# their ticket emailed.
|
|
db_name = name or ""
|
|
db_email = email or ""
|
|
|
|
db_ticket = Ticket(
|
|
id=row_id,
|
|
wallet=wallet,
|
|
event=event,
|
|
name=db_name,
|
|
email=db_email,
|
|
user_id=user_id,
|
|
registered=False,
|
|
paid=False,
|
|
reg_timestamp=now,
|
|
time=now,
|
|
extra=TicketExtra(**extra) if extra else TicketExtra(),
|
|
payment_hash=payment_hash,
|
|
)
|
|
await db.insert("events.ticket", db_ticket)
|
|
|
|
return Ticket(
|
|
id=row_id,
|
|
wallet=wallet,
|
|
event=event,
|
|
name=name,
|
|
email=email,
|
|
user_id=user_id,
|
|
registered=False,
|
|
paid=False,
|
|
reg_timestamp=now,
|
|
time=now,
|
|
extra=TicketExtra(**extra) if extra else TicketExtra(),
|
|
payment_hash=payment_hash,
|
|
)
|
|
|
|
|
|
async def update_ticket(ticket: Ticket) -> Ticket:
|
|
ticket_dict = ticket.dict()
|
|
if ticket_dict.get("name") is None:
|
|
ticket_dict["name"] = ""
|
|
if ticket_dict.get("email") is None:
|
|
ticket_dict["email"] = ""
|
|
await db.update("events.ticket", Ticket(**ticket_dict))
|
|
return ticket
|
|
|
|
|
|
async def get_tickets_by_payment_hash(payment_hash: str) -> list[Ticket]:
|
|
"""All ticket rows sharing the given LNbits invoice payment_hash.
|
|
|
|
For a single-ticket purchase returns one row (legacy invariant
|
|
`id == payment_hash` still holds). For a multi-ticket purchase
|
|
returns the N rows created with shared `payment_hash` but
|
|
distinct `id`s — each attendee's scannable QR.
|
|
"""
|
|
rows = await db.fetchall(
|
|
"SELECT * FROM events.ticket WHERE payment_hash = :ph",
|
|
{"ph": payment_hash},
|
|
)
|
|
return [Ticket(**_parse_ticket_row(row)) for row in rows]
|
|
|
|
|
|
async def get_ticket(payment_hash: str) -> Ticket | None:
|
|
row = await db.fetchone(
|
|
"SELECT * FROM events.ticket WHERE id = :id",
|
|
{"id": payment_hash},
|
|
)
|
|
if not row:
|
|
return None
|
|
return Ticket(**_parse_ticket_row(row))
|
|
|
|
|
|
async def get_tickets(wallet_ids: str | list[str]) -> list[Ticket]:
|
|
if isinstance(wallet_ids, str):
|
|
wallet_ids = [wallet_ids]
|
|
q = ",".join([f"'{wallet_id}'" for wallet_id in wallet_ids])
|
|
rows = await db.fetchall(f"SELECT * FROM events.ticket WHERE wallet IN ({q})")
|
|
return [Ticket(**_parse_ticket_row(row)) for row in rows]
|
|
|
|
|
|
async def get_tickets_by_event(event_id: str) -> list[Ticket]:
|
|
"""All ticket rows for the given calendar event id."""
|
|
rows = await db.fetchall(
|
|
"SELECT * FROM events.ticket WHERE event = :event_id",
|
|
{"event_id": event_id},
|
|
)
|
|
return [Ticket(**_parse_ticket_row(row)) for row in rows]
|
|
|
|
|
|
async def get_tickets_by_user_id(user_id: str) -> list[Ticket]:
|
|
"""All tickets owned by the given LNbits user_id."""
|
|
rows = await db.fetchall(
|
|
"SELECT * FROM events.ticket WHERE user_id = :user_id ORDER BY time DESC",
|
|
{"user_id": user_id},
|
|
)
|
|
return [Ticket(**_parse_ticket_row(row)) for row in rows]
|
|
|
|
|
|
async def get_tickets_paginated(
|
|
wallet_ids: str | list[str], filters: Filters[TicketFilters] | None = None
|
|
) -> Page[Ticket]:
|
|
if isinstance(wallet_ids, str):
|
|
wallet_ids = [wallet_ids]
|
|
|
|
wallet_where = []
|
|
values = {}
|
|
for idx, wallet_id in enumerate(wallet_ids):
|
|
key = f"wallet_id_{idx}"
|
|
wallet_where.append(f":{key}")
|
|
values[key] = wallet_id
|
|
|
|
where = [f"wallet IN ({', '.join(wallet_where)})", "paid = true"]
|
|
|
|
return await db.fetch_page(
|
|
"SELECT * FROM events.ticket",
|
|
where=where,
|
|
values=values,
|
|
filters=filters,
|
|
model=Ticket,
|
|
table_name="events.ticket",
|
|
)
|
|
|
|
|
|
async def delete_ticket(payment_hash: str) -> None:
|
|
await db.execute("DELETE FROM events.ticket WHERE id = :id", {"id": payment_hash})
|
|
|
|
|
|
async def delete_event_tickets(event_id: str) -> None:
|
|
await db.execute(
|
|
"DELETE FROM events.ticket WHERE event = :event", {"event": event_id}
|
|
)
|
|
|
|
|
|
async def purge_unpaid_tickets(event_id: str) -> None:
|
|
time_diff = datetime.now() - timedelta(hours=24)
|
|
await db.execute(
|
|
f"""
|
|
DELETE FROM events.ticket WHERE event = :event AND paid = false
|
|
AND time < {db.timestamp_placeholder("time")}
|
|
""",
|
|
{"time": time_diff.timestamp(), "event": event_id},
|
|
)
|
|
|
|
|
|
async def create_event(data: CreateEvent) -> Event:
|
|
event_id = urlsafe_short_hash()
|
|
# Default end_date to start_date and closing_date to end_date when omitted.
|
|
if not data.event_end_date:
|
|
data.event_end_date = data.event_start_date
|
|
if not data.closing_date:
|
|
data.closing_date = data.event_end_date
|
|
event = Event(id=event_id, time=datetime.now(timezone.utc), **data.dict())
|
|
event = cast(Event, sync_event_ticket_waves(event))
|
|
await db.insert("events.events", event)
|
|
return event
|
|
|
|
|
|
async def update_event(event: Event) -> Event:
|
|
event = cast(Event, sync_event_ticket_waves(event))
|
|
await db.update("events.events", event)
|
|
return event
|
|
|
|
|
|
async def get_event(event_id: str) -> Event | None:
|
|
event = await db.fetchone(
|
|
"SELECT * FROM events.events WHERE id = :id",
|
|
{"id": event_id},
|
|
Event,
|
|
)
|
|
return cast(Event, sync_event_ticket_waves(event)) if event else None
|
|
|
|
|
|
async def get_events(wallet_ids: str | list[str]) -> list[Event]:
|
|
if isinstance(wallet_ids, str):
|
|
wallet_ids = [wallet_ids]
|
|
q = ",".join([f"'{wallet_id}'" for wallet_id in wallet_ids])
|
|
events = await db.fetchall(
|
|
f"SELECT * FROM events.events WHERE wallet IN ({q})",
|
|
model=Event,
|
|
)
|
|
return [cast(Event, sync_event_ticket_waves(event)) for event in events]
|
|
|
|
|
|
async def get_all_events() -> list[Event]:
|
|
"""All events, no wallet filter. Admin-only callers."""
|
|
events = await db.fetchall(
|
|
"SELECT * FROM events.events ORDER BY time DESC",
|
|
model=Event,
|
|
)
|
|
# Wave-sync on read, exactly as upstream's `get_event` / `get_events`
|
|
# do. These four getters are fork-only, so upstream's v1.6.8 diff
|
|
# never reached them — and two of them publish: `get_all_events`
|
|
# backs /republish-all and `get_events_pending_republish` drives the
|
|
# #55 sweep. Without this they would emit the stale roll-up rather
|
|
# than the current per-wave figures.
|
|
return [cast(Event, sync_event_ticket_waves(e)) for e in events]
|
|
|
|
|
|
async def get_public_events() -> list[Event]:
|
|
"""Approved, non-canceled events for the public listing."""
|
|
events = await db.fetchall(
|
|
"""
|
|
SELECT * FROM events.events
|
|
WHERE status = 'approved' AND canceled = FALSE
|
|
ORDER BY event_start_date ASC
|
|
""",
|
|
model=Event,
|
|
)
|
|
# Wave-sync on read, exactly as upstream's `get_event` / `get_events`
|
|
# do. These four getters are fork-only, so upstream's v1.6.8 diff
|
|
# never reached them — and two of them publish: `get_all_events`
|
|
# backs /republish-all and `get_events_pending_republish` drives the
|
|
# #55 sweep. Without this they would emit the stale roll-up rather
|
|
# than the current per-wave figures.
|
|
return [cast(Event, sync_event_ticket_waves(e)) for e in events]
|
|
|
|
|
|
async def get_pending_events() -> list[Event]:
|
|
"""Proposed events awaiting admin approval."""
|
|
events = await db.fetchall(
|
|
"SELECT * FROM events.events WHERE status = 'proposed' ORDER BY time DESC",
|
|
model=Event,
|
|
)
|
|
# Wave-sync on read, exactly as upstream's `get_event` / `get_events`
|
|
# do. These four getters are fork-only, so upstream's v1.6.8 diff
|
|
# never reached them — and two of them publish: `get_all_events`
|
|
# backs /republish-all and `get_events_pending_republish` drives the
|
|
# #55 sweep. Without this they would emit the stale roll-up rather
|
|
# than the current per-wave figures.
|
|
return [cast(Event, sync_event_ticket_waves(e)) for e in events]
|
|
|
|
|
|
async def flag_wave_transitions() -> int:
|
|
"""Flag events whose advertised ticket wave has moved on.
|
|
|
|
Every republish this extension performs is *sale*-driven. A wave
|
|
boundary is a *date* boundary, so when early bird closes at midnight
|
|
nothing fires and the relay keeps serving the closed wave's price until
|
|
the next ticket happens to sell (aiolabs/events#61).
|
|
|
|
Rather than add a scheduler and a second publish path, this compares the
|
|
wave a row would advertise now against the one its last successful
|
|
publish did (`nostr_published_wave_id`) and sets `nostr_publish_pending`
|
|
on a mismatch — handing the work to the existing reconciliation sweep,
|
|
which already retries, survives restarts and logs.
|
|
|
|
Rows with NULL `nostr_published_wave_id` are skipped: that means "never
|
|
published, or published before the column existed", which is no evidence
|
|
of drift. Flagging them would republish the whole table on first boot
|
|
after the upgrade.
|
|
|
|
Returns the number of rows newly flagged.
|
|
"""
|
|
events = await db.fetchall(
|
|
"""
|
|
SELECT * FROM events.events
|
|
WHERE nostr_published_wave_id IS NOT NULL
|
|
AND nostr_publish_pending = FALSE
|
|
AND canceled = FALSE
|
|
AND status = 'approved'
|
|
""",
|
|
model=Event,
|
|
)
|
|
flagged = 0
|
|
for event in events:
|
|
event = cast(Event, sync_event_ticket_waves(event))
|
|
if advertised_wave_key(event) == event.nostr_published_wave_id:
|
|
continue
|
|
event.nostr_publish_pending = True
|
|
await update_event(event)
|
|
flagged += 1
|
|
return flagged
|
|
|
|
|
|
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.
|
|
"""
|
|
events = await db.fetchall(
|
|
"""
|
|
SELECT * FROM events.events
|
|
WHERE nostr_publish_pending = TRUE
|
|
ORDER BY time ASC
|
|
""",
|
|
model=Event,
|
|
)
|
|
# Wave-sync on read, exactly as upstream's `get_event` / `get_events`
|
|
# do. These four getters are fork-only, so upstream's v1.6.8 diff
|
|
# never reached them — and two of them publish: `get_all_events`
|
|
# backs /republish-all and `get_events_pending_republish` drives the
|
|
# #55 sweep. Without this they would emit the stale roll-up rather
|
|
# than the current per-wave figures.
|
|
return [cast(Event, sync_event_ticket_waves(e)) for e in events]
|
|
|
|
|
|
async def get_settings() -> EventsSettings:
|
|
"""Singleton settings row, seeded by m010."""
|
|
row = await db.fetchone("SELECT * FROM events.settings WHERE id = 1")
|
|
if row:
|
|
return EventsSettings(**dict(row))
|
|
return EventsSettings()
|
|
|
|
|
|
async def update_settings(settings: EventsSettings) -> EventsSettings:
|
|
await db.execute(
|
|
"UPDATE events.settings SET auto_approve = :auto_approve WHERE id = 1",
|
|
{"auto_approve": settings.auto_approve},
|
|
)
|
|
return settings
|
|
|
|
|
|
async def delete_event(event_id: str) -> None:
|
|
await db.execute("DELETE FROM events.events WHERE id = :id", {"id": event_id})
|
|
|
|
|
|
async def get_event_tickets(event_id: str) -> list[Ticket]:
|
|
rows = await db.fetchall(
|
|
"SELECT * FROM events.ticket WHERE event = :event",
|
|
{"event": event_id},
|
|
)
|
|
return [Ticket(**_parse_ticket_row(row)) for row in rows]
|