"""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 {"type": "ritual_start"} → (begins a ritual attempt) {"type": "ritual_step", "step": } → (on the final step) ritual_complete {"type": "judgment", "verdict": "trust" | "banish" | "test" | "cross_over"} → judgment_result 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 import judgment 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 ( RITUAL_SUCCESS_ESSENCE, 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 # Workstream B (character-depth-ghost-log spec): ritual progress for the # *current* entity — reset whenever a fresh presence is summoned (see # _handle_summon) or a new ritual_start arrives. ritual_steps: int = 0 ritual_completed: bool = False ritual_success: bool = False # Per-session RNG for `tell` frames — unseeded (a session's tells should # vary run to run), but persistent across calls so the draw sequence # isn't restarted on every single message. tell_rng: random.Random = field(default_factory=random.Random) # Active-session registry (spec: ESP32 sensor node, Workstream K): maps a # user to their live SeanceState so hardware ingestion (a plain HTTP # request, not a WS connection) can find "does this user have a séance # open right now" and push a hardware anomaly into it. Single # most-recent-session mapping — if a user somehow has two `/ws/session` # tabs open concurrently, the newer one wins the registry slot; that's a # reasonable v1 (see spec's Contract section) since a user's attention is # realistically in one tab at a time. Plain in-process dict, same reasoning # as everywhere else in this file: no external store needed at this scale, # doesn't need to survive a restart. _active_sessions: dict[uuid.UUID, "SeanceState"] = {} def register_active_session(user_id: uuid.UUID, state: "SeanceState") -> None: _active_sessions[user_id] = state def unregister_active_session(user_id: uuid.UUID, state: "SeanceState") -> None: # Only remove if it's still *this* state — guards against a rare # overlap where an older session's disconnect cleanup runs after a # newer session for the same user has already registered, which would # otherwise wipe out the newer (still-live) registry entry. if _active_sessions.get(user_id) is state: del _active_sessions[user_id] def get_active_session(user_id: uuid.UUID) -> "SeanceState | None": return _active_sessions.get(user_id) 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(), # Workstream B: hidden ground truth, kept on the server-side # SeanceState.entity dict for the ritual/judgment/tell handlers to # read (state.entity["traits"]) — see `_public_entity` below for why # this never reaches the wire directly. "traits": entity.traits, } def _public_entity(entity: dict) -> dict: """The entity payload actually sent to the client in the `entity` frame — everything `serialize_entity` produces *except* `traits`. Hidden traits must never leak outside `ritual_complete` on success (the contract's decoupling requirement, echoed in frontend/src/lib/evilMeter.ts's comments); `frontend/src/lib/types.ts`'s `SpiritEntity` type correspondingly has no `traits` field.""" return {key: value for key, value in entity.items() if key != "traits"} 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. An at-peace entity (Workstream B: a spirit correctly helped to cross over) is excluded from the match — it stays in the Codex forever but can't be re-contacted. If its signature is what this session's anomaly pattern hashes to, a *new* entity is minted instead. `Entity.signature` is unique, so the new entity can't reuse the exact same string while the retired row still holds it — it gets a salted variant of the same base signature instead. """ 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, Entity.at_peace.is_(False)) ) is_new = entity is None if is_new: mint_signature = signature retired = await db.scalar(select(Entity).where(Entity.signature == signature)) if retired is not None: mint_signature = f"{signature}:{uuid.uuid4().hex[:8]}" profile = await spirit_service.mint_profile( mint_signature, channel, state.anomalies, state.language ) discoverer = await db.get(User, state.user_id) favor = discoverer.favor if discoverer is not None else 0.0 entity = Entity( name=await _unique_entity_name(db, profile["name"]), epithet=profile["epithet"], persona=profile["persona"], rarity_tier=profile["rarity"], signature=mint_signature, voice_profile=profile["voice"], visual_profile=profile["visual"], sample_quotes=profile["quotes"], # Workstream B: signature-seeded traits, nudged by the # discovering user's favor (app.judgment.apply_favor_bias) — # never derived from/fed into the persona above. traits=judgment.apply_favor_bias(profile["traits"], favor), 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 (`_reward_ritual_success` / `_handle_judgment` below), which call the same `app.inventory` `roll_item_drop`/ `credit_essence` helpers and milestone essence constants 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) # A fresh presence invalidates any in-progress/completed ritual from # whatever was previously in this slot (mirrors the frontend reducer's # 'entity' case in state/seance.tsx, which resets its own ritual/ # judgment UI state the same way). state.ritual_steps = 0 state.ritual_completed = False state.ritual_success = False await state.send_queue.put( {"type": "entity", "entity": _public_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) # Workstream B: `tell` frames piggyback on the existing anomaly/reply # handling rather than running their own timer — a fragment (ambient, # frequent) rolls a lower chance than a direct reply (deliberate, a seeker # just asked something), so tells feel like they're punctuating engagement # rather than firing on a fixed clock. TELL_CHANCE_ON_FRAGMENT = 0.2 TELL_CHANCE_ON_REPLY = 0.35 async def _maybe_tell(state: SeanceState, chance: float) -> None: if state.entity is None: return if state.tell_rng.random() >= chance: return text = judgment.generate_tell(state.entity.get("traits", {}), state.tell_rng) await state.send_queue.put({"type": "tell", "text": text}) 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) await _maybe_tell(state, TELL_CHANCE_ON_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) await _maybe_tell(state, TELL_CHANCE_ON_REPLY) 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}) # --- Workstream B: ritual + judgment (character-depth-ghost-log spec) ------ # How many `ritual_step` frames complete one attempt — matches # frontend/src/lib/ritual.ts's RITUAL_TOTAL_STEPS (the 4-rune "align / # breathe / trace / lock" sequence). The frontend owns the exact step count # per the spec ("implementer's call"); this just has to agree with it. RITUAL_STEPS_REQUIRED = 4 async def _handle_ritual_start(state: SeanceState) -> None: if state.entity is None: return # no presence to focus on — frontend already gates the button state.ritual_steps = 0 state.ritual_completed = False state.ritual_success = False async def _reward_ritual_success(state: SeanceState) -> None: """Mirrors `_reward_summon`'s essence-credit + item-drop pattern for the ritual milestone trigger point.""" item = None async with session_maker() as db: user = await db.get(User, state.user_id) if user is None: return credit_essence(user, RITUAL_SUCCESS_ESSENCE) item = roll_item_drop("ritual") 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_ritual_step(state: SeanceState, message: dict) -> None: if state.entity is None or state.ritual_completed: return if not isinstance(message.get("step"), int): return state.ritual_steps += 1 if state.ritual_steps < RITUAL_STEPS_REQUIRED: return traits = state.entity.get("traits", {}) success = judgment.roll_ritual_success(traits) state.ritual_completed = True state.ritual_success = success revealed = dict(traits) if success else None await state.send_queue.put( {"type": "ritual_complete", "success": success, "revealed": revealed} ) if success: await _reward_ritual_success(state) async def _handle_judgment(state: SeanceState, message: dict) -> None: if state.entity is None: return verdict = message.get("verdict") if verdict not in judgment.VERDICTS: return traits = state.entity.get("traits", {}) outcome = judgment.judge_verdict( verdict, traits, ritual_completed=state.ritual_completed, ritual_success=state.ritual_success, ) item = None # Skip the DB round-trip entirely when there's nothing to persist (e.g. # `test` without a completed ritual, or a resisted cross_over) — the # contract's "no crash, just no effect" for those cases. if outcome.favor_delta or outcome.essence_delta or outcome.consequence in ( "reward", "crossed_over", ): async with session_maker() as db: user = await db.get(User, state.user_id) if user is not None: if outcome.favor_delta: user.favor = judgment.clamp_favor(user.favor + outcome.favor_delta) if outcome.essence_delta: credit_essence(user, outcome.essence_delta) if outcome.consequence == "crossed_over": entity_row = await db.get(Entity, uuid.UUID(state.entity["id"])) if entity_row is not None: entity_row.at_peace = True if outcome.consequence in ("reward", "crossed_over"): item = roll_item_drop("judgment") 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() await state.send_queue.put( { "type": "judgment_result", "correct": outcome.correct, "favor_delta": outcome.favor_delta, "essence_delta": outcome.essence_delta, "at_peace": outcome.at_peace, "consequence": outcome.consequence, } ) if item is not None: await state.send_queue.put({"type": "item_drop", "item": item}) @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), ) register_active_session(user_id, state) 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"))) elif msg_type == "ritual_start": await _handle_ritual_start(state) elif msg_type == "ritual_step": await _handle_ritual_step(state, message) elif msg_type == "judgment": await _handle_judgment(state, message) except WebSocketDisconnect: pass finally: unregister_active_session(user_id, state) 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