Files
plm-lottery/app/tx/confirmation.py
T

83 lines
3.2 KiB
Python
Raw Normal View History

import asyncio
import logging
from collections.abc import Awaitable, Callable
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
from app.db.models import PendingTransaction
from app.electrum.client import ElectrumClient
from app.rounds.events import broadcaster
logger = logging.getLogger(__name__)
_POLL_INTERVAL_SECONDS = 10
ConfirmationHandler = Callable[[AsyncSession, PendingTransaction], Awaitable[None]]
_handlers: dict[str, ConfirmationHandler] = {}
def register_handler(kind: str, handler: ConfirmationHandler) -> None:
"""Domain modules (bets, rounds, withdrawals) register here so this generic
poller can notify them when one of their outgoing txs gets its 1st
confirmation, without this module importing them directly."""
_handlers[kind] = handler
async def poll_once(session_factory: async_sessionmaker, client: ElectrumClient) -> int:
async with session_factory() as session:
# Plain columns, not entities: nothing then outlives the session, so this
# can't break if expire_on_commit is ever turned on (B-21).
candidates = (
await session.execute(
select(
PendingTransaction.id, PendingTransaction.current_txid, PendingTransaction.kind
).where(PendingTransaction.status == "pending")
)
).all()
confirmed = 0
for pending_id, txid, kind in candidates:
try:
tx = await client.get_transaction(txid, verbose=True)
except Exception:
# One unresolvable txid must not stop the others: a tx the server no
# longer knows (dropped from the mempool, replaced) used to abort the
# whole pass, so nothing confirmed again until an operator intervened
# (B-03). Abandoning such a row is app/tx/reconcile.py's job, not ours.
logger.warning("could not check pending_transaction %s (txid %s)", pending_id, txid, exc_info=True)
continue
if not tx or tx.get("confirmations", 0) < 1:
continue
async with session_factory() as session:
row = await session.get(PendingTransaction, pending_id)
if row is None or row.status != "pending":
continue
row.status = "confirmed"
handler = _handlers.get(kind)
if handler is not None:
await handler(session, row)
await session.commit()
broadcaster.publish() # a bet/withdrawal/payout just confirmed — balance and/or round state changed
confirmed += 1
return confirmed
class ConfirmationPoller:
def __init__(self, session_factory: async_sessionmaker, get_client: Callable[[], ElectrumClient | None]):
self._session_factory = session_factory
self._get_client = get_client
async def run(self) -> None:
while True:
client = self._get_client()
if client is not None:
try:
await poll_once(self._session_factory, client)
except asyncio.CancelledError:
raise
except Exception:
logger.exception("confirmation poll failed")
await asyncio.sleep(_POLL_INTERVAL_SECONDS)