Files
qtalker---/backend/app/ws.py
Indiana ff68379772 feat: unlocks, inventory items, sigils, drops, and essence (Workstream C)
Implements the backend REST surface and WS wiring for
docs/superpowers/specs/2026-07-23-character-depth-ghost-log-design.md's
Workstream C:

- New models: UnlockRecord (unlocks), InventoryItem (inventory_items),
  Sigil (sigils) — brand-new tables, picked up by main.py's existing
  create_all.
- New app/inventory.py: unlock price table, item drop table/odds,
  essence economy constants, sigil design validation, and an atomic
  (row-locked) purchase_unlock() that guards against double-spend races.
- New app/routes/inventory.py: GET unlocks/items/sigils, POST sigils
  (validates the placeholder {points, rune} shape, points capped at 12),
  POST unlocks/{unlock_key} (402 on insufficient essence, 404 on unknown
  key, idempotent re-buy).
- GET /auth/me now includes unlocks: list[str] and essence: int.
- ws.py: wires essence trickle + item_drop rolls into the one trigger
  point that exists in this worktree today (_handle_summon, covering
  every successful summon plus high-rarity summons); the other two
  contract trigger points (correct judgment, successful ritual) belong
  to Workstream B's not-yet-landed ritual/judgment WS handlers, which
  should call app.inventory's same helpers once they land.
- User.essence: int added (Workstream B owns this column per the spec;
  added here per orchestrator instruction so this workstream is
  independently testable — merge controller reconciles the duplicate
  edit).

Also fast-forwarded this worktree's branch onto master (it had fallen
behind several commits) so the files this workstream depends on
(shop.py, ws.py, entities.py, etc.) were actually present to build
against.

Tests: 109 passed (drop-roll statistical sanity with seeded RNG,
inventory/sigil CRUD, purchase success/insufficient-funds/idempotency/
unknown-key paths, /auth/me shape, ws summon-trickle and item-drop
wiring).

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-24 03:02:20 +00:00

534 lines
20 KiB
Python

"""The séance channel: one WebSocket per contact session.
Protocol (client → server):
{"type": "ping"} → {"type": "pong"}
{"type": "set_mode", "mode": MODE} → {"type": "mode", ...}
{"type": "language", "language": "en"|"es"}
{"type": "summon"} → {"type": "entity", ...}
{"type": "anomaly", "source": SRC, ...} → {"type": "utterance", ...}
{"type": "question", "text": "..."} → reply_start / reply_token* / reply_end
{"type": "passive", "enabled": bool} → ambient wire loop on/off
All server → client frames flow through a single sender task so concurrent
producers (ambient loop, reply streaming, TTS callbacks) never interleave on
the wire.
"""
import asyncio
import contextlib
import random
import time
import uuid
from dataclasses import dataclass, field
from datetime import datetime, timezone
from pathlib import Path
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
from sqlalchemy import select
from app.config import settings
from app.db import async_session_maker as _default_session_maker
from app.deps import SESSION_COOKIE_NAME
from app.entities import fallback_signature, signature_from_anomalies
from app.inventory import SUMMON_ESSENCE_TRICKLE, credit_essence, roll_item_drop, summon_drop_trigger
from app.llm.service import SpiritBusyError, spirit_service
from app.models.auth_session import AuthSession, hash_token
from app.models.contact_session import ContactSession
from app.models.entity import Entity
from app.models.entity_sighting import EntitySighting
from app.models.event import Event
from app.models.inventory_item import InventoryItem
from app.models.user import User
from app.possession import compute_stability
from app.rate_limit import RateLimiter, resolve_client_ip
from app.telemetry import detect_wire_spike, sample_network
from app.tts.piper import synthesize_spirit_voice
from app.tts.voices import pick_voice
router = APIRouter()
# Alias so tests can swap in the NullPool test session maker.
session_maker = _default_session_maker
MODES = {"wire", "evp", "radio", "ouija", "emf"}
# Per-user limiters for every LLM-triggering message type (spec §5).
fragment_limiter = RateLimiter(max_requests=30, window_seconds=60)
question_limiter = RateLimiter(max_requests=6, window_seconds=60)
summon_limiter = RateLimiter(max_requests=4, window_seconds=60)
# Per-IP limiters for the same trigger points (spec §5). Ollama is a shared,
# single-instance, CPU-only resource — per-account limits alone don't stop
# one compromised/scripted account from hammering it across many source IPs,
# nor do they protect the resource from many accounts sharing one IP. IP
# thresholds are looser than the per-account ones since a single IP can
# legitimately host multiple accounts on a shared network.
fragment_ip_limiter = RateLimiter(max_requests=60, window_seconds=60)
question_ip_limiter = RateLimiter(max_requests=12, window_seconds=60)
summon_ip_limiter = RateLimiter(max_requests=8, window_seconds=60)
AUDIO_DIR = Path(settings.data_dir) / "audio"
def _audio_dir() -> Path:
AUDIO_DIR.mkdir(parents=True, exist_ok=True)
return AUDIO_DIR
@dataclass
class SeanceState:
user_id: uuid.UUID
session_id: uuid.UUID
client_ip: str
send_queue: asyncio.Queue = field(default_factory=asyncio.Queue)
mode: str = "unknown"
language: str = "en"
entity: dict | None = None
anomalies: list[dict] = field(default_factory=list)
history: list[dict] = field(default_factory=list)
ambient_task: asyncio.Task | None = None
wire_jitter_history: list[float] = field(default_factory=list)
last_wire_anomaly_at: float = 0.0
def serialize_entity(entity: Entity) -> dict:
return {
"id": str(entity.id),
"name": entity.name,
"epithet": entity.epithet,
"persona": entity.persona,
"rarity": entity.rarity_tier,
"voice": entity.voice_profile,
"visual": entity.visual_profile,
"quotes": entity.sample_quotes,
"contact_count": entity.contact_count,
"discovered_at": entity.discovered_at.isoformat(),
}
def _client_ip(websocket: WebSocket) -> str:
host = websocket.client.host if websocket.client else None
return resolve_client_ip(websocket.headers, host)
async def _authenticate(websocket: WebSocket) -> uuid.UUID | None:
raw_token = websocket.cookies.get(SESSION_COOKIE_NAME)
if raw_token is None:
return None
token_hash = hash_token(raw_token)
async with session_maker() as db:
result = await db.execute(
select(AuthSession).where(AuthSession.token_hash == token_hash)
)
session = result.scalar_one_or_none()
if session is None or session.expires_at < datetime.now(timezone.utc):
return None
return session.user_id
async def _sender(state: SeanceState, websocket: WebSocket) -> None:
"""The only task allowed to write to the socket."""
try:
while True:
message = await state.send_queue.get()
await websocket.send_json(message)
except asyncio.CancelledError:
pass
async def _record_event(
session_id: uuid.UUID,
kind: str,
text: str | None = None,
payload: dict | None = None,
) -> uuid.UUID:
async with session_maker() as db:
event = Event(session_id=session_id, kind=kind, text=text, payload=payload)
db.add(event)
await db.commit()
await db.refresh(event)
return event.id
async def _speak(
state: SeanceState, kind: str, text: str, instability: float = 0.0
) -> None:
"""Persist an utterance, push its text immediately, synthesize audio in
the background, and push the audio URL when the effects chain finishes."""
event_id = await _record_event(
state.session_id, "utterance", text=text, payload={"kind": kind}
)
await state.send_queue.put(
{
"type": "utterance",
"id": str(event_id),
"kind": kind,
"text": text,
"entity": state.entity["name"] if state.entity else None,
}
)
async def _synth() -> None:
try:
voice_profile = state.entity.get("voice", {}) if state.entity else {}
voice = pick_voice(voice_profile.get("voice_id"), state.language)
wav = await synthesize_spirit_voice(
text, voice, voice_profile, instability=instability
)
filename = f"{event_id}.wav"
_audio_dir().joinpath(filename).write_bytes(wav)
async with session_maker() as db:
event = await db.get(Event, event_id)
if event is not None:
event.audio_path = filename
await db.commit()
await state.send_queue.put(
{"type": "audio", "id": str(event_id), "url": f"/audio/{filename}"}
)
except Exception:
# TTS is texture, not content — the words already reached the
# seeker. Never let a synth failure kill the séance.
pass
asyncio.create_task(_synth())
_ROMAN = ["II", "III", "IV", "V", "VI", "VII", "VIII", "IX"]
async def _unique_entity_name(db, base_name: str) -> str:
name = base_name
i = 0
while await db.scalar(select(Entity).where(Entity.name == name)) is not None:
name = f"{base_name} {_ROMAN[i]}" if i < len(_ROMAN) else f"{base_name} {i + 2}"
i += 1
return name
async def _summon(state: SeanceState, channel: str) -> tuple[Entity, bool]:
"""Match this session's signature against the Codex, or mint a new entity."""
signature = signature_from_anomalies(state.anomalies) or fallback_signature(
str(state.session_id)
)
async with session_maker() as db:
entity = await db.scalar(select(Entity).where(Entity.signature == signature))
is_new = entity is None
if is_new:
profile = await spirit_service.mint_profile(
signature, channel, state.anomalies, state.language
)
entity = Entity(
name=await _unique_entity_name(db, profile["name"]),
epithet=profile["epithet"],
persona=profile["persona"],
rarity_tier=profile["rarity"],
signature=signature,
voice_profile=profile["voice"],
visual_profile=profile["visual"],
sample_quotes=profile["quotes"],
discovered_by=state.user_id,
contact_count=1,
)
db.add(entity)
else:
entity.contact_count += 1
await db.flush()
session = await db.get(ContactSession, state.session_id)
if session is not None:
session.entity_id = entity.id
db.add(
EntitySighting(
entity_id=entity.id, session_id=state.session_id, user_id=state.user_id
)
)
await db.commit()
await db.refresh(entity)
return entity, is_new
async def _reward_summon(state: SeanceState) -> None:
"""Workstream C's essence-trickle + item-drop trigger points that exist
in this handler today: a small essence trickle for every successful
summon (any mode), and — only when the summoned entity is high-rarity —
a roll for an item drop. The other two contract trigger points ("after a
correct judgment, a successful ritual") belong to Workstream B's
ritual/judgment WS handlers, which don't exist in this codebase yet;
`app.inventory` exposes the same `roll_item_drop`/`credit_essence`
helpers (plus the milestone essence constants) for those handlers to
call once they land, so the drop table and essence economy stay in one
place instead of being duplicated."""
assert state.entity is not None
rarity = state.entity.get("rarity", "common")
async with session_maker() as db:
user = await db.get(User, state.user_id)
if user is None:
return
credit_essence(user, SUMMON_ESSENCE_TRICKLE)
item = None
trigger = summon_drop_trigger(rarity)
if trigger is not None:
item = roll_item_drop(trigger)
if item is not None:
db.add(
InventoryItem(
user_id=state.user_id,
item_type=item["item_type"],
item_key=item["item_key"],
payload=item["payload"],
)
)
await db.commit()
if item is not None:
await state.send_queue.put({"type": "item_drop", "item": item})
async def _handle_summon(state: SeanceState) -> None:
# Short-circuits: an account already over its own cap never gets far
# enough to spend from the IP budget too.
if not (
summon_limiter.allow(str(state.user_id))
and summon_ip_limiter.allow(state.client_ip)
):
await state.send_queue.put(
{
"type": "error",
"code": "rate_limited",
"message": "The veil is crowded. The spirits need a moment before another summoning.",
}
)
return
await state.send_queue.put({"type": "status", "state": "summoning"})
entity, is_new = await _summon(state, state.mode if state.mode != "unknown" else "ouija")
state.entity = serialize_entity(entity)
await state.send_queue.put(
{"type": "entity", "entity": state.entity, "is_new": is_new}
)
await _reward_summon(state)
greeting = random.choice(state.entity["quotes"]) if state.entity["quotes"] else "I am here."
await _speak(state, "greeting", greeting)
async def _handle_anomaly(state: SeanceState, message: dict) -> None:
anomaly = {
"source": str(message.get("source", "unknown"))[:16],
"frequency": message.get("frequency"),
"magnitude": message.get("magnitude"),
}
state.anomalies.append(anomaly)
state.anomalies = state.anomalies[-64:]
await _record_event(state.session_id, "anomaly", payload=anomaly)
await state.send_queue.put({"type": "anomaly_ack", "count": len(state.anomalies)})
if state.entity is None:
if signature_from_anomalies(state.anomalies) is not None:
# The signal has enough structure — something announces itself.
await _handle_summon(state)
else:
await state.send_queue.put({"type": "status", "state": "attuning"})
return
if not (
fragment_limiter.allow(str(state.user_id))
and fragment_ip_limiter.allow(state.client_ip)
):
return # anomalies during a crowded veil just pass unheard
try:
fragment = await spirit_service.fragment(
anomaly["source"], anomaly, state.language
)
except SpiritBusyError:
return
await _speak(state, "fragment", fragment)
async def _handle_question(state: SeanceState, text: str) -> None:
if not (
question_limiter.allow(str(state.user_id))
and question_ip_limiter.allow(state.client_ip)
):
await state.send_queue.put(
{
"type": "error",
"code": "rate_limited",
"message": "The spirit is spent. Give it a moment to gather itself.",
}
)
return
if state.entity is None:
await _handle_summon(state)
assert state.entity is not None
text = text.strip()[:500]
await _record_event(state.session_id, "question", text=text)
await state.send_queue.put({"type": "status", "state": "gathering"})
last_magnitude = state.anomalies[-1].get("magnitude") if state.anomalies else 50
if last_magnitude is None:
last_magnitude = 50
stability = compute_stability(state.entity["rarity"], float(last_magnitude))
await state.send_queue.put({"type": "reply_start", "stability": stability})
reply_parts: list[str] = []
try:
async for token in spirit_service.chat_stream(
state.entity, text, state.history, state.language
):
reply_parts.append(token)
await state.send_queue.put({"type": "reply_token", "token": token})
except SpiritBusyError:
await state.send_queue.put(
{
"type": "error",
"code": "veil_crowded",
"message": "Too many seekers press against the veil. The spirit withdraws.",
}
)
await state.send_queue.put({"type": "reply_end", "text": ""})
return
reply = "".join(reply_parts).strip()
reply_id = await _record_event(state.session_id, "reply", text=reply)
state.history.append({"role": "user", "text": text})
state.history.append({"role": "spirit", "text": reply})
state.history = state.history[-8:]
await state.send_queue.put({"type": "reply_end", "id": str(reply_id), "text": reply})
if reply:
await _speak(state, "reply", reply, instability=1 - stability)
async def _ambient_loop(state: SeanceState) -> None:
"""The Wire Ghost's pulse: telemetry every few seconds. Spikes in the
jitter baseline become first-class anomaly events (driving fragments and
the client's data views); quieter ticks may earn an ambient whisper."""
try:
while True:
await asyncio.sleep(random.uniform(3, 5))
try:
sample = await sample_network(period_s=1.0)
except Exception:
continue
await state.send_queue.put({"type": "telemetry", **sample.as_dict()})
jitter = sample.jitter_bytes_per_s
ratio = detect_wire_spike(state.wire_jitter_history, jitter)
state.wire_jitter_history = (state.wire_jitter_history + [jitter])[-60:]
now = time.monotonic()
if ratio is not None and now - state.last_wire_anomaly_at > 20:
state.last_wire_anomaly_at = now
anomaly = {
"source": "wire",
"frequency": round(jitter, 1),
"magnitude": round(ratio * 10, 1),
}
# Tell the client first so its matrix/graph views light up,
# then run the anomaly through the usual séance pipeline.
await state.send_queue.put({"type": "anomaly", **anomaly})
await _handle_anomaly(state, anomaly)
continue
if not spirit_service.ambient_ready():
continue
whisper = await spirit_service.wire_whisper(sample.as_dict(), state.language)
await _speak(state, "ambient", whisper)
except asyncio.CancelledError:
pass
async def _handle_passive(state: SeanceState, enabled: bool) -> None:
if enabled and (state.ambient_task is None or state.ambient_task.done()):
state.ambient_task = asyncio.create_task(_ambient_loop(state))
await state.send_queue.put({"type": "passive", "enabled": True})
elif not enabled and state.ambient_task is not None:
state.ambient_task.cancel()
state.ambient_task = None
await state.send_queue.put({"type": "passive", "enabled": False})
@router.websocket("/ws/session")
async def session_socket(websocket: WebSocket) -> None:
user_id = await _authenticate(websocket)
if user_id is None:
await websocket.close(code=4401)
return
await websocket.accept()
async with session_maker() as db:
contact_session = ContactSession(user_id=user_id)
db.add(contact_session)
await db.commit()
await db.refresh(contact_session)
state = SeanceState(
user_id=user_id,
session_id=contact_session.id,
client_ip=_client_ip(websocket),
)
sender = asyncio.create_task(_sender(state, websocket))
await state.send_queue.put({"type": "session", "id": str(contact_session.id)})
try:
while True:
message = await websocket.receive_json()
msg_type = message.get("type")
if msg_type == "ping":
await state.send_queue.put({"type": "pong"})
elif msg_type == "set_mode" and message.get("mode") in MODES:
state.mode = message["mode"]
async with session_maker() as db:
session = await db.get(ContactSession, state.session_id)
if session is not None:
session.mode = state.mode
await db.commit()
await state.send_queue.put({"type": "mode", "mode": state.mode})
elif msg_type == "language" and message.get("language") in ("en", "es"):
state.language = message["language"]
async with session_maker() as db:
session = await db.get(ContactSession, state.session_id)
if session is not None:
session.language = state.language
await db.commit()
elif msg_type == "summon":
await _handle_summon(state)
elif msg_type == "anomaly":
await _handle_anomaly(state, message)
elif msg_type == "question" and isinstance(message.get("text"), str):
await _handle_question(state, message["text"])
elif msg_type == "passive":
await _handle_passive(state, bool(message.get("enabled")))
except WebSocketDisconnect:
pass
finally:
if state.ambient_task is not None:
state.ambient_task.cancel()
sender.cancel()
with contextlib.suppress(asyncio.CancelledError):
await sender
# ASGI servers may cancel the handler task as soon as the socket
# closes, killing any await here mid-flight — so the session-close
# write runs detached, surviving the handler's own teardown.
asyncio.create_task(_close_session(contact_session.id))
async def _close_session(session_id: uuid.UUID) -> None:
try:
async with session_maker() as db:
session_to_close = await db.get(ContactSession, session_id)
if session_to_close is not None:
session_to_close.ended_at = datetime.now(timezone.utc)
await db.commit()
except Exception:
pass