A dropped connection used to hang the whole platform permanently, and three defects composed to do it (BUGS.md B-01): The read loop's death was invisible. When the socket closed, _read_loop broke out and finished, but _run_once was blocked on gather() over two notification consumers waiting on queues nobody would ever fill again — it never returned and never raised, so the reconnect-with-backoff logic was unreachable. client.wait_closed() now resolves when the loop ends for any reason, and _run_once races it against the consumers and a keepalive with asyncio.wait(FIRST_COMPLETED). Nothing had a timeout. request() registered a future, wrote to a half-closed socket (drain() often doesn't raise) and awaited a reply that would never come. That hung a POST /bets *while holding the per-user lock*, and could stop the confirmation poller for good. Every request is now bounded at 15s, and a timeout tears the connection down rather than leaving a server that owes us a reply in rotation. There was no keepalive, so on a quiet instance the normal way this connection dies is an idle-timeout drop by the server (~10 minutes for many). A server.ping every 60s makes that observable within a minute. listener.client is also cleared before reconnecting, so callers stop treating a dead connection as live. On top of the finding, the listener now rotates over a list of servers: ELECTRUM_FALLBACK_SERVERS holds comma-separated host:port[:notls] extras, tried after the primary. Everything the platform does goes through this one connection — deposit credits, broadcasts, confirmations, the chain tip the draw waits on — which made a single hardcoded server its biggest point of failure. A failed or dropped session moves to the next server immediately and only sleeps on the backoff once every server has had a turn, so one dead server costs one attempt instead of an outage, while a genuinely offline network still backs off. A malformed entry fails at startup, not during the outage when the fallback is what you need. Also fixes B-19: header handling refuses a height below the current tip and applies height and hex together, since _wait_for_next_block waits for tip_height > tip_at_close (a regression silently added a block to the draw's wait) and that hex is the draw's entropy source, so a mismatched pair would be worse than a stale one. Verified in the live deployment: the log shows the endpoint list, then "Electrum connected to santantonio.sytes.net:50002", and the connection holds. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
164 lines
6.2 KiB
Python
164 lines
6.2 KiB
Python
import pytest
|
|
from cryptography.fernet import Fernet
|
|
from httpx import ASGITransport, AsyncClient
|
|
|
|
from app.config import settings
|
|
|
|
|
|
@pytest.fixture
|
|
async def client(monkeypatch, tmp_path):
|
|
monkeypatch.setattr(settings, "database_url", f"sqlite+aiosqlite:///{tmp_path}/test.db")
|
|
monkeypatch.setattr(settings, "jwt_secret", "test-jwt-secret")
|
|
monkeypatch.setattr(settings, "xprv_encryption_key", Fernet.generate_key().decode())
|
|
monkeypatch.setattr(settings, "master_key_path", str(tmp_path / "master.xprv.enc"))
|
|
|
|
import app.wallet.hd as hd
|
|
|
|
hd._account_key = None
|
|
hd.generate_master_key()
|
|
|
|
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
|
|
|
|
from app.db import base as db_base
|
|
|
|
import app.db.models # noqa: F401
|
|
|
|
db_base.engine = create_async_engine(settings.database_url)
|
|
db_base.AsyncSessionLocal = async_sessionmaker(db_base.engine, expire_on_commit=False)
|
|
|
|
from app.db import session as db_session
|
|
|
|
db_session.AsyncSessionLocal = db_base.AsyncSessionLocal
|
|
|
|
async with db_base.engine.begin() as conn:
|
|
await conn.run_sync(db_base.Base.metadata.create_all)
|
|
|
|
from fastapi import FastAPI
|
|
|
|
from app.api.routes.rounds import router as rounds_router
|
|
from app.auth.routes import router as auth_router
|
|
from app.electrum.listener import ElectrumListener
|
|
|
|
app = FastAPI()
|
|
app.include_router(auth_router)
|
|
app.include_router(rounds_router)
|
|
app.state.electrum_listener = ElectrumListener(lambda endpoint: None, db_base.AsyncSessionLocal)
|
|
|
|
transport = ASGITransport(app=app)
|
|
async with AsyncClient(transport=transport, base_url="http://test") as ac:
|
|
yield ac, db_base.AsyncSessionLocal
|
|
|
|
await db_base.engine.dispose()
|
|
|
|
|
|
async def _register(ac, username):
|
|
resp = await ac.post("/auth/register", json={"username": username, "password": "hunter2hunter"})
|
|
assert resp.status_code == 201
|
|
data = resp.json()
|
|
return data["access_token"], data["user_id"] if "user_id" in data else None
|
|
|
|
|
|
async def test_user_played_true_only_for_participants(client):
|
|
ac, session_factory = client
|
|
from app.db.models import Round, RoundConfig, RoundParticipant, User
|
|
|
|
player_token, _ = await _register(ac, "player")
|
|
spectator_token, _ = await _register(ac, "spectator")
|
|
|
|
async with session_factory() as session:
|
|
from sqlalchemy import select
|
|
|
|
session.add(RoundConfig(fee_address="pool-fee-address"))
|
|
player = (await session.scalars(select(User).where(User.username == "player"))).one()
|
|
round_ = Round(status="paying_out", winner_user_id=player.id, winner_amount_sats=123)
|
|
session.add(round_)
|
|
await session.flush()
|
|
session.add(
|
|
RoundParticipant(
|
|
round_id=round_.id,
|
|
user_id=player.id,
|
|
bet_amount_sats=1_000_000_000,
|
|
bet_txid="a" * 64,
|
|
)
|
|
)
|
|
await session.commit()
|
|
|
|
resp = await ac.get("/rounds/current", headers={"Authorization": f"Bearer {player_token}"})
|
|
assert resp.status_code == 200
|
|
assert resp.json()["user_played"] is True
|
|
|
|
resp = await ac.get("/rounds/current", headers={"Authorization": f"Bearer {spectator_token}"})
|
|
assert resp.status_code == 200
|
|
assert resp.json()["user_played"] is False
|
|
|
|
resp = await ac.get("/rounds/current") # no auth at all — logged-out chain-only view
|
|
assert resp.status_code == 200
|
|
assert resp.json()["user_played"] is False
|
|
|
|
|
|
async def test_jackpot_comes_from_the_participants_actual_bets(client):
|
|
"""B-11: the jackpot was participant_count * the *current* bet_amount_sats, which
|
|
overstated it (each stored bet is already net of that bet's network fee) and
|
|
silently rewrote the advertised jackpot of a round in progress whenever an
|
|
operator edited the bet amount."""
|
|
from sqlalchemy import select
|
|
|
|
from app.db.models import Round, RoundConfig, RoundParticipant
|
|
|
|
ac, session_factory = client
|
|
|
|
async with session_factory() as session:
|
|
session.add(RoundConfig(fee_address="", bet_amount_sats=1_000_000_000))
|
|
session.add(Round(id=50, status="open"))
|
|
await session.flush()
|
|
# Two bets that actually paid 999_800_000 each (fee deducted), not 1_000_000_000.
|
|
session.add(
|
|
RoundParticipant(round_id=50, user_id=1, bet_amount_sats=999_800_000, bet_txid="a", status="confirmed")
|
|
)
|
|
session.add(
|
|
RoundParticipant(round_id=50, user_id=2, bet_amount_sats=999_800_000, bet_txid="b", status="confirmed")
|
|
)
|
|
await session.commit()
|
|
|
|
body = (await ac.get("/rounds/current")).json()
|
|
assert body["participant_count"] == 2
|
|
assert body["jackpot_sats"] == (999_800_000 * 2) * 70 // 100
|
|
|
|
# Changing the configured bet amount must not move a running round's jackpot.
|
|
async with session_factory() as session:
|
|
config = (await session.scalars(select(RoundConfig))).one()
|
|
config.bet_amount_sats = 5_000_000_000
|
|
await session.commit()
|
|
|
|
body = (await ac.get("/rounds/current")).json()
|
|
assert body["jackpot_sats"] == (999_800_000 * 2) * 70 // 100
|
|
|
|
|
|
async def test_unhandled_errors_use_the_structured_detail_shape(client):
|
|
"""B-24: the catch-all handler answered with a bare-string `detail`, while
|
|
app/api/errors.py documents detail as {"code", "message", "params"}. Clients then
|
|
had to special-case exactly the responses they understand least."""
|
|
from fastapi import FastAPI
|
|
from httpx import ASGITransport, AsyncClient
|
|
|
|
from app.main import log_unhandled_exception
|
|
|
|
app = FastAPI()
|
|
app.add_exception_handler(Exception, log_unhandled_exception)
|
|
|
|
@app.get("/boom")
|
|
async def boom():
|
|
raise RuntimeError("secret internal detail")
|
|
|
|
transport = ASGITransport(app=app, raise_app_exceptions=False)
|
|
async with AsyncClient(transport=transport, base_url="http://test") as ac:
|
|
resp = await ac.get("/boom")
|
|
|
|
assert resp.status_code == 500
|
|
detail = resp.json()["detail"]
|
|
assert detail["code"] == "internal_error"
|
|
assert detail["message"] == "internal server error"
|
|
assert detail["params"] == {}
|
|
# The exception text belongs in logs/app.log, never in the response body.
|
|
assert "secret internal detail" not in resp.text
|