events/crud.py
Padreug 15c2276e57
Some checks failed
lint.yml / feat(nostr): publish the active ticket wave, not the roll-up (pull_request) Failing after 0s
feat(nostr): publish the active ticket wave, not the roll-up
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.
2026-09-28 22:28:09 +02:00

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]