"""WS integration tests for the DURABILITY of the reward guards in app.ws. `test_ws_passage.py` and `test_ws_reward_replay.py` closed three replay faucets — the ritual milestone, a verdict, and a Passage beat — but they closed them with sets held on the in-memory `SeanceState`, and every one of those tests replays the exploit down a SINGLE websocket. That is exactly the shape of the remaining hole: the guards were per-connection, so a client that dropped the socket and reconnected, or simply opened a second one alongside the first, got a fresh empty guard and could be paid all over again for the same spirit. Income stayed rate-bounded by `summon_limiter` (per user id, across connections) and unbounded in total. These tests reach past one connection: * walk a rite to completion, DISCONNECT, reconnect on a new socket with the same account, land back on the same spirit, re-walk it — the ledger must not grow; * two sockets open SIMULTANEOUSLY for the same account, both claiming the same award — paid exactly once; * the same race driven concurrently rather than interleaved, straight at `_claim_award`, which is where the UNIQUE constraint on award_claims actually decides the winner; * and the anti-overshoot direction in both places: a genuinely fresh presence still pays in full across a reconnect, and a corrected verdict still pays. LANDING ON THE SAME SPIRIT ACROSS CONNECTIONS is the load-bearing detail. A spirit is matched by the *channel* signature, and with no anomalies at all that signature falls back to the CONTACT SESSION id — which is per-connection, so a bare reconnect+summon mints a brand new spirit and would prove nothing. So these tests feed the same anomaly stream a browser feeds (which is also what auto-summons), giving both connections the same channel, and pin the `answers` draw so the familiar presence reliably answers rather than 72% of the time. """ import asyncio import uuid import pytest from sqlalchemy import select import app.judgment as judgment_module import app.ws from app.entities import MIN_ANOMALIES_FOR_SIGNATURE, fallback_profile from app.inventory import ( CORRECT_JUDGMENT_ESSENCE, RITUAL_SUCCESS_ESSENCE, SUMMON_ESSENCE_TRICKLE, ) from app.models.award_claim import AwardClaim from app.models.inventory_item import InventoryItem from app.models.user import User from app.rate_limit import RateLimiter class FakeSpiritService: async def mint_profile( self, signature, channel, anomalies, language="en", entropy=None, sky=None ): return fallback_profile(signature) async def fragment(self, source, anomaly, language="en"): return "listen" async def wire_whisper(self, telemetry, language="en"): return "the wire hums" def chat_stream(self, entity, question, history, language="en"): async def gen(): yield "here." return gen() def ambient_ready(self): return False async def _fake_synth(text, voice, profile, instability=0.0): return b"RIFFfake wav bytes" @pytest.fixture(autouse=True) def _fake_spirits(monkeypatch): monkeypatch.setattr(app.ws, "spirit_service", FakeSpiritService()) monkeypatch.setattr(app.ws, "synthesize_spirit_voice", _fake_synth) # Module-level limiters are shared singletons accumulating real hit counts # across the whole session (precedent: test_ws_reward_replay.py). These # tests deliberately summon more than the real 4/60s cap allows, because # the point is what happens on the SECOND connection — a rate limit there # would mask the durability question instead of answering it. for name in ( "summon_limiter", "summon_ip_limiter", "question_limiter", "question_ip_limiter", "fragment_limiter", "fragment_ip_limiter", "ritual_limiter", "ritual_ip_limiter", "judgment_limiter", "judgment_ip_limiter", "passage_limiter", "passage_ip_limiter", ): monkeypatch.setattr(app.ws, name, RateLimiter(max_requests=10_000, window_seconds=60)) def _read_until(ws, msg_type, max_frames=120, **match): for _ in range(max_frames): frame = ws.receive_json() if frame.get("type") != msg_type: continue if all(frame.get(key) == value for key, value in match.items()): return frame raise AssertionError(f"never saw frame of type {msg_type!r} matching {match!r}") def _login(sync_client, username): sync_client.post("/auth/register", json={"username": username, "password": "spookyspooky"}) sync_client.post("/auth/login", json={"username": username, "password": "spookyspooky"}) return sync_client.cookies.get("qm_session") def _user_id(sync_client, token): return uuid.UUID( sync_client.get("/auth/me", headers={"cookie": f"qm_session={token}"}).json()["id"] ) def _ws_connect(sync_client, token): return sync_client.websocket_connect( "/ws/session", headers={"cookie": f"qm_session={token}"} ) def _settle(ws): """Force a full round trip so any post-frame DB write has landed. The connection's message loop is sequential, so a pong is a hard guarantee that the previous handler returned — not a poll-and-hope.""" ws.send_json({"type": "ping"}) _read_until(ws, "pong") # One fixed anomaly pattern = one fixed channel signature, shared by every # connection in a test. `signature_from_anomalies` buckets frequency by digit # count and sorts magnitudes, so identical frames give an identical channel # regardless of which contact session sent them. _CHANNEL = [ {"type": "anomaly", "source": "emf", "frequency": 121.5, "magnitude": 30.0}, {"type": "anomaly", "source": "emf", "frequency": 122.5, "magnitude": 31.0}, {"type": "anomaly", "source": "emf", "frequency": 123.5, "magnitude": 32.0}, ] assert len(_CHANNEL) >= MIN_ANOMALIES_FOR_SIGNATURE def _familiar_answers(monkeypatch): """Pin the summon draw so a known channel ALWAYS returns its familiar presence. Untouched, this is a genuine draw against RETURN_CHANCE — roughly 7 times in 10. These tests need the two connections to be looking at the SAME spirit for the question ("can it be paid twice?") to mean anything, so a draw of 0.0 is used, which is below any possible `return_chance`. Every other draw stays real. """ real = app.ws.veil_float def selective(entropy, context, *args, **kwargs): if context == "answers": return 0.0 return real(entropy, context, *args, **kwargs) monkeypatch.setattr(app.ws, "veil_float", selective) def _open_channel(ws): """Drive the channel a browser drives: enough anomalies to fingerprint it, which is itself what triggers the summon (see `_handle_anomaly`).""" _read_until(ws, "session") for frame in _CHANNEL: ws.send_json(frame) entity = _read_until(ws, "entity") _settle(ws) return entity["entity"]["id"] def _run_ritual(ws, steps=4): """One full attempt: the rewind button, then every step.""" ws.send_json({"type": "ritual_start"}) for i in range(1, steps + 1): ws.send_json({"type": "ritual_step", "step": i}) result = _read_until(ws, "ritual_complete") _settle(ws) return result def _judge(ws, verdict): ws.send_json({"type": "judgment", "verdict": verdict}) result = _read_until(ws, "judgment_result") _settle(ws) return result async def _ledger(db_session, user_id): db_session.expire_all() user = await db_session.get(User, user_id) items = ( await db_session.execute( select(InventoryItem).where(InventoryItem.user_id == user_id) ) ).scalars().all() return user.essence, user.favor, len(items) def _always_wins(monkeypatch): monkeypatch.setattr( app.ws.judgment, "roll_ritual_success", lambda traits, rng=None: True ) def _verdict_outcomes(monkeypatch, table): monkeypatch.setattr( app.ws.judgment, "judge_verdict", lambda verdict, traits, **kw: table[verdict] ) _CORRECT_TRUST = judgment_module.JudgmentOutcome( True, judgment_module.FAVOR_CORRECT_TRUST, CORRECT_JUDGMENT_ESSENCE, False, "reward" ) _WRONG_TRUST = judgment_module.JudgmentOutcome( False, judgment_module.FAVOR_WRONG_TRUST, 0, False, "escalation" ) _CORRECT_BANISH = judgment_module.JudgmentOutcome( True, judgment_module.FAVOR_CORRECT_BANISH, CORRECT_JUDGMENT_ESSENCE, False, "reward" ) # --- reconnecting ---------------------------------------------------------- # How many times each reconnect test re-opens the socket. Named because the # expected summon trickle is derived from it, and a bare 5 in two places that # must agree is exactly how that arithmetic drifts. RECONNECTS = 5 @pytest.mark.asyncio async def test_reconnecting_does_not_reopen_the_ritual_purse( sync_client, db_session, monkeypatch ): """The exploit the in-memory guard could not see, executed literally. Complete the rite, DROP THE SOCKET, reconnect on a fresh websocket with the same account, tune back to the same channel so the same spirit answers, and complete the rite again — five times over. The seeker must finish with exactly one milestone's worth of essence, because there was only ever one spirit. `ritual_paid` was a bool on `SeanceState`, which dies with the connection, so before this every reconnect paid RITUAL_SUCCESS_ESSENCE again and rolled another item. """ _always_wins(monkeypatch) _familiar_answers(monkeypatch) token = _login(sync_client, "ritual-reconnector") user_id = _user_id(sync_client, token) with _ws_connect(sync_client, token) as ws: first_entity = _open_channel(ws) baseline, _, _ = await _ledger(db_session, user_id) assert _run_ritual(ws)["success"] is True after_first, _, items_after_first = await _ledger(db_session, user_id) for _ in range(RECONNECTS): with _ws_connect(sync_client, token) as ws: entity_id = _open_channel(ws) assert entity_id == first_entity, ( "the reconnect landed on a different spirit — this test cannot " "say anything about paying the same one twice" ) assert _run_ritual(ws)["success"] is True, "the replay did not even complete" after_reconnects, _, items_after_reconnects = await _ledger(db_session, user_id) assert after_first - baseline == RITUAL_SUCCESS_ESSENCE, ( "the rite paid nothing — the test proves nothing" ) # Each reconnect re-opens the channel, and opening a channel is a summon, # which pays SUMMON_ESSENCE_TRICKLE every time by design (app/inventory.py:42 # — "every successful summon, any mode"; it is bounded by summon_limiter at # 4/60s, not by any once-per-presence rule). That trickle is NOT what this # test is about, so it is expected explicitly rather than folded into the # comparison: what must not repeat is the 15-essence ritual MILESTONE. expected_trickle = RECONNECTS * SUMMON_ESSENCE_TRICKLE assert after_reconnects == after_first + expected_trickle, ( f"reconnecting minted {after_reconnects - after_first - expected_trickle} extra " f"essence beyond the {expected_trickle} of expected summon trickle, across " f"{RECONNECTS} fresh websockets; the guard dies with the connection" ) assert items_after_reconnects == items_after_first, ( f"reconnecting rolled {items_after_reconnects - items_after_first} extra item " "drops off a single spirit" ) @pytest.mark.asyncio async def test_reconnecting_does_not_reopen_the_judgment_purse( sync_client, db_session, monkeypatch ): """Same hole, the verdict road. `judged_verdicts` was a set on `SeanceState`; every reconnect handed the client an empty one, so the same verdict on the same spirit paid essence AND favor AND an item roll again, with favor walking toward its +1.0 ceiling (which then biases every future mint).""" _verdict_outcomes(monkeypatch, {"trust": _CORRECT_TRUST}) _familiar_answers(monkeypatch) token = _login(sync_client, "judgment-reconnector") user_id = _user_id(sync_client, token) with _ws_connect(sync_client, token) as ws: first_entity = _open_channel(ws) baseline_essence, baseline_favor, _ = await _ledger(db_session, user_id) _judge(ws, "trust") essence_one, favor_one, items_one = await _ledger(db_session, user_id) replays = [] for _ in range(RECONNECTS): with _ws_connect(sync_client, token) as ws: assert _open_channel(ws) == first_entity replays.append(_judge(ws, "trust")) essence_end, favor_end, items_end = await _ledger(db_session, user_id) assert essence_one - baseline_essence == CORRECT_JUDGMENT_ESSENCE assert favor_one - baseline_favor == pytest.approx(judgment_module.FAVOR_CORRECT_TRUST) # As in the ritual test: re-opening the channel is a summon, and the # trickle is paid per summon by design. The VERDICT is what must not repeat. # Favor, checked just below, has no trickle — so it must not move at all. expected_trickle = RECONNECTS * SUMMON_ESSENCE_TRICKLE assert essence_end == essence_one + expected_trickle, ( f"reconnecting minted {essence_end - essence_one - expected_trickle} extra " f"essence beyond the {expected_trickle} of expected summon trickle, across " f"{RECONNECTS} fresh websockets; the verdict guard dies with the connection" ) assert favor_end == pytest.approx(favor_one), ( f"reconnecting moved favor by a further {favor_end - favor_one}; a client " "that reconnects in a loop can pin favor at the +1.0 ceiling" ) assert items_end == items_one for frame in replays: assert frame["essence_delta"] == 0 and frame["favor_delta"] == 0, ( "a verdict replayed on a fresh connection advertised essence/favor " "it did not pay" ) @pytest.mark.asyncio async def test_reconnecting_does_not_reopen_the_passage_purse( sync_client, db_session, monkeypatch ): """Same hole, the Passage road. `passage_paid` was a set on `SeanceState`, so every reconnect re-credited the whole ladder of beats.""" real = app.ws.veil_float def selective(entropy, context, *args, **kwargs): # No collapses, no lies — the totals must depend on the guard, not the # draw — and the familiar presence always answers. if context.startswith("passage:"): return 1.0 if context == "answers": return 0.0 return real(entropy, context, *args, **kwargs) monkeypatch.setattr(app.ws, "veil_float", selective) # Every beat but `release`: crossing latches `at_peace`, which correctly # takes the spirit out of reach of any later summon. beats = len(app.ws.passage.LAYERS) - 1 token = _login(sync_client, "passage-reconnector") user_id = _user_id(sync_client, token) def _walk(ws): for _ in range(beats): ws.send_json({"type": "passage_layer"}) _read_until(ws, "passage_result") _settle(ws) with _ws_connect(sync_client, token) as ws: first_entity = _open_channel(ws) baseline, _, _ = await _ledger(db_session, user_id) ws.send_json({"type": "passage_start"}) _settle(ws) _walk(ws) after_first, _, _ = await _ledger(db_session, user_id) for _ in range(RECONNECTS): with _ws_connect(sync_client, token) as ws: assert _open_channel(ws) == first_entity ws.send_json({"type": "passage_start"}) _settle(ws) _walk(ws) after_reconnects, _, _ = await _ledger(db_session, user_id) assert after_first - baseline > 0, "the rite paid nothing — the test proves nothing" # Re-opening the channel summons, and the trickle is per-summon by design; # the BEATS are what must not pay twice. expected_trickle = RECONNECTS * SUMMON_ESSENCE_TRICKLE assert after_reconnects == after_first + expected_trickle, ( f"reconnecting minted {after_reconnects - after_first - expected_trickle} extra " f"essence beyond the {expected_trickle} of expected summon trickle, across " f"{RECONNECTS} fresh websockets; the passage is an unbounded faucet again" ) # --- two sockets at once --------------------------------------------------- @pytest.mark.asyncio async def test_two_simultaneous_connections_pay_the_ritual_once( sync_client, db_session, monkeypatch ): """Not a reconnect: both sockets are OPEN AT THE SAME TIME. This is the cheaper half of the exploit — no disconnect needed, just a second tab. Each connection has its own `SeanceState` and therefore had its own empty guard, so the two of them paid the same spirit's ritual twice while every replay test in the suite watched a single socket. """ _always_wins(monkeypatch) _familiar_answers(monkeypatch) token = _login(sync_client, "ritual-two-tabs") user_id = _user_id(sync_client, token) with _ws_connect(sync_client, token) as first, _ws_connect(sync_client, token) as second: entity_a = _open_channel(first) entity_b = _open_channel(second) assert entity_a == entity_b, ( "the two connections landed on different spirits — this test cannot " "say anything about paying the same one twice" ) baseline, _, items_baseline = await _ledger(db_session, user_id) # Interleaved, so neither connection's rite is finished before the # other's begins. first.send_json({"type": "ritual_start"}) second.send_json({"type": "ritual_start"}) for i in range(1, 5): first.send_json({"type": "ritual_step", "step": i}) second.send_json({"type": "ritual_step", "step": i}) assert _read_until(first, "ritual_complete")["success"] is True assert _read_until(second, "ritual_complete")["success"] is True _settle(first) _settle(second) essence_end, _, items_end = await _ledger(db_session, user_id) assert essence_end - baseline == RITUAL_SUCCESS_ESSENCE, ( f"two simultaneous connections paid {essence_end - baseline} for one " f"spirit's ritual instead of {RITUAL_SUCCESS_ESSENCE}" ) assert items_end - items_baseline <= 1, ( "two simultaneous connections each rolled the ritual's item drop" ) @pytest.mark.asyncio async def test_two_simultaneous_connections_pay_a_verdict_once( sync_client, db_session, monkeypatch ): _verdict_outcomes(monkeypatch, {"trust": _CORRECT_TRUST}) _familiar_answers(monkeypatch) token = _login(sync_client, "judgment-two-tabs") user_id = _user_id(sync_client, token) with _ws_connect(sync_client, token) as first, _ws_connect(sync_client, token) as second: assert _open_channel(first) == _open_channel(second) baseline_essence, baseline_favor, _ = await _ledger(db_session, user_id) first.send_json({"type": "judgment", "verdict": "trust"}) second.send_json({"type": "judgment", "verdict": "trust"}) frames = [ _read_until(first, "judgment_result"), _read_until(second, "judgment_result"), ] _settle(first) _settle(second) essence_end, favor_end, _ = await _ledger(db_session, user_id) assert essence_end - baseline_essence == CORRECT_JUDGMENT_ESSENCE, ( f"two simultaneous connections paid {essence_end - baseline_essence} for one " f"verdict instead of {CORRECT_JUDGMENT_ESSENCE}" ) assert favor_end - baseline_favor == pytest.approx( judgment_module.FAVOR_CORRECT_TRUST ), "two simultaneous connections moved favor twice for one verdict" # And the frames must be honest about it: exactly one of them paid. assert sum(f["essence_delta"] for f in frames) == essence_end - baseline_essence @pytest.mark.asyncio async def test_a_concurrent_claim_race_resolves_to_one_winner(sync_client, db_session): """Straight at the constraint, with real concurrency. The two-tabs tests above interleave frames but the server still handles them one at a time, so they prove the guard is DURABLE without ever forcing the two claims to be in flight together. This one does: twelve coroutines claim the same (user, entity, award) at once via `asyncio.gather`. Exactly one may win. If `_claim_award` were a check-then-insert, several would read "unclaimed" and all of them would pay; if it did not catch the IntegrityError, the losers would raise instead of returning False and would take down their whole séance with a 500. """ token = _login(sync_client, "claim-racer") user_id = _user_id(sync_client, token) with _ws_connect(sync_client, token) as ws: entity_id = _open_channel(ws) state = app.ws.SeanceState( user_id=user_id, session_id=uuid.uuid4(), client_ip="127.0.0.1" ) async def claim(): # A private state per racer: sharing one would let the in-memory fast # path answer some of them, which is exactly what must NOT decide this. racer = app.ws.SeanceState( user_id=user_id, session_id=state.session_id, client_ip="127.0.0.1" ) return await app.ws._claim_award(racer, entity_id, "ritual") results = await asyncio.gather(*(claim() for _ in range(12))) assert sum(results) == 1, ( f"{sum(results)} of 12 concurrent claims won the same award; the " "UNIQUE constraint is not deciding the race" ) db_session.expire_all() rows = ( await db_session.execute( select(AwardClaim).where( AwardClaim.user_id == user_id, AwardClaim.entity_id == uuid.UUID(entity_id), AwardClaim.award_key == "ritual", ) ) ).scalars().all() assert len(rows) == 1, f"the race left {len(rows)} claim rows for one award" # --- the anti-overshoot direction ------------------------------------------ @pytest.mark.asyncio async def test_a_fresh_presence_after_a_reconnect_still_pays_in_full( sync_client, db_session, monkeypatch ): """The guard must not overshoot across connections either. A durable claim that keyed on the user alone — or on anything coarser than the presence — would silently make every spirit after the first worthless for the rest of that account's life, which is a far worse bug than the faucet. Reconnect onto a DIFFERENT channel, meet a different spirit, and the rite must pay in full again. """ _always_wins(monkeypatch) token = _login(sync_client, "fresh-after-reconnect") user_id = _user_id(sync_client, token) with _ws_connect(sync_client, token) as ws: first_entity = _open_channel(ws) baseline, _, _ = await _ledger(db_session, user_id) _run_ritual(ws) after_first, _, _ = await _ledger(db_session, user_id) # A different anomaly pattern is a different channel, so a different # presence answers — the honest version of what the replay tests fake. with _ws_connect(sync_client, token) as ws: _read_until(ws, "session") for i in range(MIN_ANOMALIES_FOR_SIGNATURE): ws.send_json( {"type": "anomaly", "source": "radio", "frequency": 9000.0 + i, "magnitude": 77.0 + i} ) second_entity = _read_until(ws, "entity")["entity"]["id"] _settle(ws) assert second_entity != first_entity, ( "the second channel produced the same spirit — this test is about a " "genuinely new presence" ) _run_ritual(ws) after_second, _, _ = await _ledger(db_session, user_id) first_rite = after_first - baseline assert first_rite == RITUAL_SUCCESS_ESSENCE # The second connection's summon pays its own trickle too, so compare only # the amount that arrived once the rite completed. assert after_second - after_first >= RITUAL_SUCCESS_ESSENCE, ( "a fresh presence met on a new connection did not re-open the purse — " "every spirit after the first is worth nothing" ) @pytest.mark.asyncio async def test_a_correction_still_pays_on_a_second_connection( sync_client, db_session, monkeypatch ): """The claim is keyed PER VERDICT, and that must survive the move to the database. A seeker who calls it wrong on one connection and right on the next is correcting themselves, not replaying — a blanket "judged once per spirit" row would make the corrected verdict a no-op.""" _verdict_outcomes(monkeypatch, {"trust": _WRONG_TRUST, "banish": _CORRECT_BANISH}) _familiar_answers(monkeypatch) token = _login(sync_client, "corrector-reconnect") user_id = _user_id(sync_client, token) with _ws_connect(sync_client, token) as ws: first_entity = _open_channel(ws) baseline_essence, baseline_favor, _ = await _ledger(db_session, user_id) wrong = _judge(ws, "trust") with _ws_connect(sync_client, token) as ws: assert _open_channel(ws) == first_entity right = _judge(ws, "banish") essence_end, favor_end, _ = await _ledger(db_session, user_id) assert wrong["consequence"] == "escalation" assert right["consequence"] == "reward" # One reconnect here, so exactly one summon trickle rides along with the # corrected verdict's payout. assert essence_end - baseline_essence == CORRECT_JUDGMENT_ESSENCE + SUMMON_ESSENCE_TRICKLE, ( "the corrected verdict paid nothing on the second connection — the " "durable guard is too broad" ) assert favor_end - baseline_favor == pytest.approx( judgment_module.FAVOR_WRONG_TRUST + judgment_module.FAVOR_CORRECT_BANISH ) @pytest.mark.asyncio async def test_the_test_verdict_stays_freely_repeatable_across_connections( sync_client, db_session, monkeypatch ): """`test` is a diagnostic pulse that never touches the ledger, so it is deliberately exempt — it must never write a claim row, or the panel would stop answering a seeker who pulses the same spirit from a second tab.""" _familiar_answers(monkeypatch) token = _login(sync_client, "tester-reconnect") user_id = _user_id(sync_client, token) results = [] for _ in range(3): with _ws_connect(sync_client, token) as ws: entity_id = _open_channel(ws) baseline, baseline_favor, _ = await _ledger(db_session, user_id) results.append(_judge(ws, "test")) essence_end, favor_end, _ = await _ledger(db_session, user_id) assert essence_end == baseline and favor_end == baseline_favor assert all(r["consequence"] == "neutral" for r in results) db_session.expire_all() claims = ( await db_session.execute( select(AwardClaim).where( AwardClaim.user_id == user_id, AwardClaim.entity_id == uuid.UUID(entity_id), ) ) ).scalars().all() assert claims == [], "the `test` verdict wrote a claim row and is no longer free"