Add transaction broadcast, confirmation polling and RBF fee-bump

Shared per-user locking to serialize bet/withdrawal PSBT builds
(tx/locks.py), a confirmation poller for pending outgoing transactions,
and the timeout->fee-bump->rebroadcast loop used by bets, payouts and
withdrawals alike (tx/broadcast.py).

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
2026-07-21 10:25:57 +02:00
co-authored by Claude Sonnet 5
parent 107e592704
commit fc2aadbc7e
6 changed files with 515 additions and 0 deletions
+181
View File
@@ -0,0 +1,181 @@
from datetime import datetime, timedelta, timezone
import pytest
from embit import script
from embit.bip32 import HDKey
from embit.transaction import Transaction
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
from app.config import settings
from app.db.base import Base
from app.db.models import PendingTransaction, User
from app.tx.broadcast import RbfError, bump_fee, should_bump
from app.wallet.plm_network import PLM_MAINNET
from app.wallet.psbt_builder import Utxo, build_signed_transaction
def _key(seed_byte: int) -> HDKey:
root = HDKey.from_seed(bytes([seed_byte]) * 32, version=PLM_MAINNET["xprv"])
return root.derive("m/84h/746h/0h/0/0")
def test_should_bump_false_before_timeout():
pending = PendingTransaction(
kind="bet", current_txid="x", fee_rate_sat_vb=1, raw_tx_hex="00", status="pending",
broadcast_at=datetime.now(timezone.utc),
)
assert should_bump(pending, datetime.now(timezone.utc), timeout_seconds=900) is False
def test_should_bump_true_after_timeout():
pending = PendingTransaction(
kind="bet", current_txid="x", fee_rate_sat_vb=1, raw_tx_hex="00", status="pending",
broadcast_at=datetime.now(timezone.utc) - timedelta(seconds=1000),
)
assert should_bump(pending, datetime.now(timezone.utc), timeout_seconds=900) is True
def test_should_bump_false_when_not_pending():
pending = PendingTransaction(
kind="bet", current_txid="x", fee_rate_sat_vb=1, raw_tx_hex="00", status="confirmed",
broadcast_at=datetime.now(timezone.utc) - timedelta(seconds=1000),
)
assert should_bump(pending, datetime.now(timezone.utc), timeout_seconds=900) is False
class FakeClient:
def __init__(self, prevout_values: dict[str, int]):
self._prevout_values = prevout_values
self.broadcasted: list[str] = []
async def get_transaction(self, txid: str, verbose: bool = False) -> dict:
return {"vout": {0: {"value": self._prevout_values[txid] / 100_000_000}}}
async def broadcast(self, raw_tx_hex: str) -> str:
self.broadcasted.append(raw_tx_hex)
return "network-txid"
@pytest.fixture
async def session_factory(tmp_path, monkeypatch):
monkeypatch.setattr(settings, "master_key_path", str(tmp_path / "master.xprv.enc"))
monkeypatch.setattr(
settings,
"xprv_encryption_key",
__import__("cryptography.fernet", fromlist=["Fernet"]).Fernet.generate_key().decode(),
)
from app.wallet import hd
hd._account_key = None
hd.generate_master_key()
engine = create_async_engine("sqlite+aiosqlite:///:memory:")
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
yield async_sessionmaker(engine, expire_on_commit=False)
await engine.dispose()
hd._account_key = None
async def test_bump_fee_shrinks_change_and_rebroadcasts(session_factory):
from app.wallet.hd import derive_user_address, derive_user_key
signer = derive_user_key(0)
my_address = derive_user_address(0)
from_script = script.p2wpkh(signer.to_public())
to_address = script.p2wpkh(_key(99).to_public()).address(network=PLM_MAINNET)
utxo_amount = 150_000_000
utxo_txid = "11" * 32
built = build_signed_transaction(
signing_key=signer,
from_script=from_script,
utxos=[Utxo(utxo_txid, 0, utxo_amount)],
to_address=to_address,
amount_sats=10_000_000,
change_address=my_address,
fee_rate_sat_vb=1,
)
async with session_factory() as session:
user = User(username="alice", password_hash="x", derivation_index=0, address=my_address)
session.add(user)
await session.commit()
pending = PendingTransaction(
kind="bet",
user_id=user.id,
current_txid=built.txid,
fee_rate_sat_vb=1,
raw_tx_hex=built.raw_hex,
status="pending",
broadcast_at=datetime.now(timezone.utc) - timedelta(seconds=1000),
)
session.add(pending)
await session.commit()
pending_id = pending.id
client = FakeClient({utxo_txid: utxo_amount})
async with session_factory() as session:
row = await session.get(PendingTransaction, pending_id)
new_txid = await bump_fee(session, client, row)
assert client.broadcasted
assert new_txid != built.txid
new_tx = Transaction.parse(bytes.fromhex(client.broadcasted[0]))
old_tx = Transaction.parse(bytes.fromhex(built.raw_hex))
old_change = next(o.value for o in old_tx.vout if o.script_pubkey.address(network=PLM_MAINNET) == my_address)
new_change = next(o.value for o in new_tx.vout if o.script_pubkey.address(network=PLM_MAINNET) == my_address)
assert new_change < old_change # fee bump came out of the change output
async with session_factory() as session:
row = await session.get(PendingTransaction, pending_id)
assert row.current_txid == new_txid
assert row.fee_rate_sat_vb == 2
assert row.attempt_count == 2
async def test_bump_fee_raises_when_no_change_output(session_factory):
from app.wallet.hd import derive_user_address, derive_user_key
signer = derive_user_key(0)
my_address = derive_user_address(0)
from_script = script.p2wpkh(signer.to_public())
to_address = script.p2wpkh(_key(98).to_public()).address(network=PLM_MAINNET)
utxo_amount = 10_000_000 # exact amount, no change output
utxo_txid = "22" * 32
built = build_signed_transaction(
signing_key=signer,
from_script=from_script,
utxos=[Utxo(utxo_txid, 0, utxo_amount)],
to_address=to_address,
amount_sats=10_000_000,
change_address=my_address,
fee_rate_sat_vb=1,
)
async with session_factory() as session:
user = User(username="bob", password_hash="x", derivation_index=0, address=my_address)
session.add(user)
await session.commit()
pending = PendingTransaction(
kind="bet",
user_id=user.id,
current_txid=built.txid,
fee_rate_sat_vb=1,
raw_tx_hex=built.raw_hex,
status="pending",
broadcast_at=datetime.now(timezone.utc) - timedelta(seconds=1000),
)
session.add(pending)
await session.commit()
pending_id = pending.id
client = FakeClient({utxo_txid: utxo_amount})
async with session_factory() as session:
row = await session.get(PendingTransaction, pending_id)
with pytest.raises(RbfError):
await bump_fee(session, client, row)
+82
View File
@@ -0,0 +1,82 @@
import pytest
from sqlalchemy import select
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
import app.bets.confirmation # noqa: F401 (registers the "bet" handler)
import app.rounds.confirmation # noqa: F401 (registers the "payout" handler)
from app.db.base import Base
from app.db.models import PendingTransaction, Round, RoundParticipant
from app.tx.confirmation import poll_once
class FakeClient:
def __init__(self, confirmations_by_txid: dict[str, int]):
self._confirmations = confirmations_by_txid
async def get_transaction(self, txid: str, verbose: bool = False) -> dict:
return {"confirmations": self._confirmations.get(txid, 0)}
@pytest.fixture
async def session_factory():
engine = create_async_engine("sqlite+aiosqlite:///:memory:")
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
yield async_sessionmaker(engine, expire_on_commit=False)
await engine.dispose()
async def test_bet_confirmation_marks_participant_confirmed(session_factory):
async with session_factory() as session:
session.add(Round(id=1, status="open"))
session.add(
RoundParticipant(
round_id=1, user_id=1, bet_amount_sats=1_000, bet_txid="tx1", status="broadcast"
)
)
session.add(
PendingTransaction(kind="bet", round_id=1, user_id=1, current_txid="tx1", fee_rate_sat_vb=1, raw_tx_hex="00", status="pending")
)
await session.commit()
client = FakeClient({"tx1": 1})
confirmed = await poll_once(session_factory, client)
assert confirmed == 1
async with session_factory() as session:
participant = (await session.scalars(select(RoundParticipant))).one()
assert participant.status == "confirmed"
assert participant.confirmed_at is not None
pending = (await session.scalars(select(PendingTransaction))).one()
assert pending.status == "confirmed"
async def test_unconfirmed_tx_is_left_pending(session_factory):
async with session_factory() as session:
session.add(Round(id=2, status="open"))
session.add(RoundParticipant(round_id=2, user_id=1, bet_amount_sats=1_000, bet_txid="tx2", status="broadcast"))
session.add(PendingTransaction(kind="bet", round_id=2, user_id=1, current_txid="tx2", fee_rate_sat_vb=1, raw_tx_hex="00", status="pending"))
await session.commit()
client = FakeClient({"tx2": 0})
confirmed = await poll_once(session_factory, client)
assert confirmed == 0
async with session_factory() as session:
participant = (await session.scalars(select(RoundParticipant))).one()
assert participant.status == "broadcast"
async def test_payout_confirmation_closes_round(session_factory):
async with session_factory() as session:
session.add(Round(id=3, status="paying_out", payout_txid="tx3"))
session.add(PendingTransaction(kind="payout", round_id=3, current_txid="tx3", fee_rate_sat_vb=1, raw_tx_hex="00", status="pending"))
await session.commit()
client = FakeClient({"tx3": 2})
confirmed = await poll_once(session_factory, client)
assert confirmed == 1
async with session_factory() as session:
round_ = await session.get(Round, 3)
assert round_.status == "closed"