_release_failed_bet and _release_failed_withdrawal restored the balance, freed the reserved UTXOs and (for a bet) removed the participant without calling broadcaster.publish(), so every dashboard kept showing the phantom bet and the reduced balance until its next poll — while the success path and the reconciler's own abandon path both published. The two regression tests pre-open the round before subscribing: place_bet opens one itself, and that publish() would otherwise satisfy the assertion whether or not the rollback published anything.
263 lines
11 KiB
Python
263 lines
11 KiB
Python
from datetime import datetime, timedelta, timezone
|
|
|
|
import pytest
|
|
from sqlalchemy import select
|
|
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
|
|
|
|
from app.bets.service import BetError, place_bet
|
|
from app.config import settings
|
|
from app.db.base import Base
|
|
from app.db.models import AuditLog, PendingTransaction, Round, RoundConfig, RoundParticipant, User, UtxoEvent
|
|
from app.rounds.events import broadcaster
|
|
from app.rounds.service import open_new_round_if_needed
|
|
from app.wallet.hd import derive_user_address
|
|
from app.wallet.psbt_builder import MAX_TX_INPUTS
|
|
|
|
|
|
class FakeElectrumClient:
|
|
def __init__(self):
|
|
self.broadcasted: list[str] = []
|
|
|
|
async def broadcast(self, raw_tx_hex: str) -> str:
|
|
self.broadcasted.append(raw_tx_hex)
|
|
return "fake-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 _make_funded_user(session_factory, index: int, funded_sats: int) -> int:
|
|
async with session_factory() as session:
|
|
address = derive_user_address(index)
|
|
user = User(username=f"user{index}", password_hash="x", derivation_index=index, address=address)
|
|
session.add(user)
|
|
await session.commit()
|
|
session.add(
|
|
UtxoEvent(
|
|
user_id=user.id,
|
|
txid=f"{index:02x}" * 32,
|
|
vout=0,
|
|
amount_sats=funded_sats,
|
|
confirmed_height=100,
|
|
)
|
|
)
|
|
await session.commit()
|
|
return user.id
|
|
|
|
|
|
async def test_place_bet_broadcasts_and_records_participant(session_factory):
|
|
user_id = await _make_funded_user(session_factory, 0, 1_500_000_000)
|
|
client = FakeElectrumClient()
|
|
|
|
async with session_factory() as session:
|
|
user = await session.get(User, user_id)
|
|
participant = await place_bet(session, client, user)
|
|
|
|
assert client.broadcasted # a raw tx was broadcast
|
|
assert participant.status == "broadcast"
|
|
assert participant.bet_txid
|
|
|
|
async with session_factory() as session:
|
|
utxo = (await session.scalars(select(UtxoEvent).where(UtxoEvent.user_id == user_id))).one()
|
|
assert utxo.spent_txid == participant.bet_txid
|
|
|
|
pending = (await session.scalars(select(PendingTransaction))).one()
|
|
assert pending.kind == "bet"
|
|
|
|
audit_events = (await session.scalars(select(AuditLog))).all()
|
|
assert any(e.event_type == "bet_placed" for e in audit_events)
|
|
assert pending.current_txid == participant.bet_txid
|
|
|
|
|
|
async def test_place_bet_rejects_insufficient_balance(session_factory):
|
|
user_id = await _make_funded_user(session_factory, 1, 1_000_000) # below bet_amount_sats
|
|
client = FakeElectrumClient()
|
|
|
|
async with session_factory() as session:
|
|
user = await session.get(User, user_id)
|
|
with pytest.raises(BetError, match="insufficient balance"):
|
|
await place_bet(session, client, user)
|
|
|
|
|
|
async def test_place_bet_reports_a_too_fragmented_balance_distinctly(session_factory): # B-48
|
|
# 100 x 0.15 PLM = 15 PLM, plenty for a 10 PLM bet, but the 50 largest inputs
|
|
# only add up to 7.5 PLM — so the build must fail with its own code, not with
|
|
# the "you have no funds" one, and must carry the cap for the translation.
|
|
user_id = await _make_funded_user(session_factory, 20, 15_000_000)
|
|
async with session_factory() as session:
|
|
for i in range(99):
|
|
session.add(
|
|
UtxoEvent(
|
|
user_id=user_id,
|
|
txid=f"{i:064x}",
|
|
vout=0,
|
|
amount_sats=15_000_000,
|
|
confirmed_height=100,
|
|
)
|
|
)
|
|
await session.commit()
|
|
client = FakeElectrumClient()
|
|
|
|
async with session_factory() as session:
|
|
user = await session.get(User, user_id)
|
|
with pytest.raises(BetError) as excinfo:
|
|
await place_bet(session, client, user)
|
|
|
|
assert excinfo.value.code == "too_many_inputs"
|
|
assert excinfo.value.params == {"max_inputs": MAX_TX_INPUTS}
|
|
assert not client.broadcasted
|
|
|
|
|
|
async def test_place_bet_rejects_second_bet_same_round(session_factory):
|
|
user_id = await _make_funded_user(session_factory, 2, 3_000_000_000)
|
|
client = FakeElectrumClient()
|
|
|
|
async with session_factory() as session:
|
|
user = await session.get(User, user_id)
|
|
await place_bet(session, client, user)
|
|
|
|
async with session_factory() as session:
|
|
user = await session.get(User, user_id)
|
|
with pytest.raises(BetError, match="already"):
|
|
await place_bet(session, client, user)
|
|
|
|
async with session_factory() as session:
|
|
participants = (await session.scalars(select(RoundParticipant))).all()
|
|
assert len(participants) == 1
|
|
|
|
|
|
async def test_place_bet_rejects_after_timer_expires_even_if_still_open(session_factory):
|
|
"""The scheduler only flips status "open" -> "closing" on its next tick (up
|
|
to a few seconds late) — place_bet must independently refuse bets once the
|
|
round's own deadline has passed, so no new player can sneak in during that
|
|
gap (see rounds/service.round_accepts_bets)."""
|
|
user_id = await _make_funded_user(session_factory, 3, 3_000_000_000)
|
|
client = FakeElectrumClient()
|
|
|
|
async with session_factory() as session:
|
|
session.add(RoundConfig(fee_address="", round_duration_seconds=60))
|
|
round_ = await open_new_round_if_needed(session)
|
|
round_.opened_at = datetime.now(timezone.utc) - timedelta(seconds=61)
|
|
await session.commit()
|
|
|
|
async with session_factory() as session:
|
|
user = await session.get(User, user_id)
|
|
with pytest.raises(BetError, match="closing"):
|
|
await place_bet(session, client, user)
|
|
|
|
async with session_factory() as session:
|
|
participants = (await session.scalars(select(RoundParticipant))).all()
|
|
assert len(participants) == 0
|
|
round_ = (await session.scalars(select(Round))).one()
|
|
assert round_.status == "open" # scheduler hasn't ticked — status is unchanged, only the check is deadline-aware
|
|
|
|
|
|
class RejectingElectrumClient:
|
|
"""A node that refuses the transaction — fee too low, dust output, mempool
|
|
conflict, or simply an unreachable server."""
|
|
|
|
async def broadcast(self, raw_tx_hex: str) -> str:
|
|
raise RuntimeError("min relay fee not met")
|
|
|
|
|
|
async def test_failed_broadcast_leaves_nothing_behind(session_factory):
|
|
"""B-07/B-08: the broadcast used to happen before anything was written, so a
|
|
rejection left the UTXOs marked spent with no rows to explain it, and the caller
|
|
got an opaque HTTP 500. Now it's a translatable error and a full rollback."""
|
|
user_id = await _make_funded_user(session_factory, 4, 3_000_000_000)
|
|
|
|
async with session_factory() as session:
|
|
user = await session.get(User, user_id)
|
|
with pytest.raises(BetError, match="refused"):
|
|
await place_bet(session, RejectingElectrumClient(), user)
|
|
|
|
async with session_factory() as session:
|
|
utxo = (await session.scalars(select(UtxoEvent).where(UtxoEvent.user_id == user_id))).one()
|
|
assert utxo.spent_txid is None # released, so the user can bet again
|
|
assert (await session.scalars(select(RoundParticipant))).all() == []
|
|
assert (await session.scalars(select(PendingTransaction))).all() == []
|
|
user = await session.get(User, user_id)
|
|
assert user.cached_balance_sats == 3_000_000_000
|
|
events = [e.event_type for e in (await session.scalars(select(AuditLog))).all()]
|
|
assert "bet_broadcast_failed" in events
|
|
assert "bet_placed" not in events
|
|
|
|
|
|
async def test_failed_broadcast_publishes_an_sse_update(session_factory): # B-49
|
|
"""The rollback moves as much state as the successful path does, so it must ping
|
|
the dashboards the same way — otherwise the phantom bet stays on screen until the
|
|
next poll."""
|
|
user_id = await _make_funded_user(session_factory, 21, 3_000_000_000)
|
|
async with session_factory() as session:
|
|
# Open the round up front: place_bet would otherwise open it itself, and that
|
|
# publish() would satisfy the assertion below whether or not the rollback ever
|
|
# published one of its own.
|
|
await open_new_round_if_needed(session)
|
|
await session.commit()
|
|
|
|
queue = broadcaster.subscribe()
|
|
try:
|
|
while not queue.empty():
|
|
queue.get_nowait()
|
|
|
|
async with session_factory() as session:
|
|
user = await session.get(User, user_id)
|
|
with pytest.raises(BetError, match="refused"):
|
|
await place_bet(session, RejectingElectrumClient(), user)
|
|
|
|
assert not queue.empty()
|
|
finally:
|
|
broadcaster.unsubscribe(queue)
|
|
|
|
|
|
async def test_failed_broadcast_reports_the_broadcast_failed_code(session_factory):
|
|
user_id = await _make_funded_user(session_factory, 5, 3_000_000_000)
|
|
|
|
async with session_factory() as session:
|
|
user = await session.get(User, user_id)
|
|
try:
|
|
await place_bet(session, RejectingElectrumClient(), user)
|
|
assert False, "expected BetError"
|
|
except BetError as exc:
|
|
assert exc.code == "broadcast_failed"
|
|
|
|
|
|
async def test_bet_is_persisted_before_it_is_broadcast(session_factory):
|
|
"""The ordering guarantee behind B-08: by the time the network call happens, the
|
|
rows already exist, so a crash there is recoverable rather than silent."""
|
|
user_id = await _make_funded_user(session_factory, 6, 3_000_000_000)
|
|
seen: dict[str, object] = {}
|
|
|
|
class ObservingClient:
|
|
async def broadcast(self, raw_tx_hex: str) -> str:
|
|
# Read committed state from an independent session, mid-broadcast.
|
|
async with session_factory() as probe:
|
|
seen["pending"] = [
|
|
(p.kind, p.status) for p in (await probe.scalars(select(PendingTransaction))).all()
|
|
]
|
|
seen["participants"] = [
|
|
(p.status) for p in (await probe.scalars(select(RoundParticipant))).all()
|
|
]
|
|
return "network-txid"
|
|
|
|
async with session_factory() as session:
|
|
user = await session.get(User, user_id)
|
|
await place_bet(session, ObservingClient(), user)
|
|
|
|
assert seen["pending"] == [("bet", "building")]
|
|
assert seen["participants"] == ["building"]
|