Unbounded essence/item farming: every other reward trigger (summon,
question, fragment) had both a per-user and per-IP limiter, but
ritual_start and judgment had none at all — and judgment has no
"already resolved" state either. A scripted client could replay
{"type":"judgment","verdict":"cross_over"} in a tight loop and mint
CROSS_OVER_ESSENCE (25) plus a 20% item roll every iteration, forever.
Same for ritual_start -> 4x ritual_step. Added ritual/judgment limiters
in both flavors, matching the existing pattern.
WS session crash: _handle_question did `if state.entity is None:
await _handle_summon(state)` then `assert state.entity is not None`.
_handle_summon returns early *without* setting state.entity when the
seeker is rate-limited, so the assert fired unhandled — and the message
loop only catches WebSocketDisconnect, so it killed the whole connection.
Reachable with no malice: click summon a few times impatiently, then ask a
question. Now returns cleanly (the rate_limited frame was already sent).
Essence double-spend: purchase_unlock() deliberately uses SELECT ... FOR
UPDATE to serialize concurrent purchases, but the three credit_essence
call sites in ws.py did an unlocked db.get() read-modify-write. An
unlocked read doesn't block on a row lock, so a reward computed from a
pre-purchase balance could be written after the purchase committed,
silently reverting the deduction — user keeps the unlock and the essence.
All three now lock the row the same way.
Entity mint collision: _summon does a racy check-then-insert against
Entity.signature and Entity.name, both DB-unique, with no IntegrityError
handling — a concurrent mint of the same signature crashed the session.
Forceable by a user with two accounts (anomaly frequency/magnitude are
client-controlled), and plausible without malice in wire mode, where
sample_network() reads host-wide /proc/net/dev counters so two idle
sessions genuinely measure the same traffic. Now retries once, which
re-runs the match against whatever the winner committed.
Also added a unique constraint on unlocks(user_id, unlock_key) as
defense-in-depth, with an idempotent catalog-guarded migration.
221 backend tests pass.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
851 lines
34 KiB
Python
851 lines
34 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
|
|
{"type": "ritual_start"} → (begins a ritual attempt)
|
|
{"type": "ritual_step", "step": <int>} → (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 sqlalchemy.exc import IntegrityError
|
|
|
|
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)
|
|
|
|
# ritual_start/judgment don't call the LLM, but each one is an essence/favor/
|
|
# item-drop reward trigger point same as summon/question — unlike those,
|
|
# they previously had no limiter at all, which meant a scripted client could
|
|
# credit itself unbounded essence by simply replaying "judgment" (or
|
|
# ritual_start -> 4x ritual_step) in a tight loop. These bound that to the
|
|
# same modest, human-plausible cadence as the other reward triggers.
|
|
ritual_limiter = RateLimiter(max_requests=6, window_seconds=60)
|
|
judgment_limiter = RateLimiter(max_requests=10, 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)
|
|
ritual_ip_limiter = RateLimiter(max_requests=12, window_seconds=60)
|
|
judgment_ip_limiter = RateLimiter(max_requests=20, 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)
|
|
)
|
|
|
|
# Two concurrent sessions can race to mint the same signature (this is
|
|
# the whole point of the match-or-mint check above being racy across
|
|
# connections), or two mints can independently decide on the same
|
|
# "next available" name via _unique_entity_name's read-then-decide
|
|
# check — either raises IntegrityError on commit against Entity's
|
|
# unique(signature)/unique(name) constraints. Retrying re-runs the
|
|
# match against what the winning transaction just committed, so the
|
|
# loser finds and reuses that row instead of crashing the session.
|
|
last_error: IntegrityError | None = None
|
|
for _attempt in range(2):
|
|
async with session_maker() as db:
|
|
try:
|
|
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
|
|
except IntegrityError as exc:
|
|
await db.rollback()
|
|
last_error = exc
|
|
continue
|
|
|
|
assert last_error is not None
|
|
raise last_error
|
|
|
|
|
|
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:
|
|
# Locked (see purchase_unlock's docstring in app/inventory.py) so
|
|
# this read-modify-write on essence can't race a concurrent
|
|
# purchase's own locked deduction and silently clobber it.
|
|
user = await db.scalar(select(User).where(User.id == state.user_id).with_for_update())
|
|
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)
|
|
if state.entity is None:
|
|
# _handle_summon() returns without setting state.entity when
|
|
# the seeker is rate-limited (it already sent its own
|
|
# "rate_limited" error frame in that case) — bail out here
|
|
# instead of asserting, which would raise uncaught and crash
|
|
# this session's whole WS message loop (only WebSocketDisconnect
|
|
# is caught around it in session_socket()).
|
|
return
|
|
|
|
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
|
|
if not (
|
|
ritual_limiter.allow(str(state.user_id))
|
|
and ritual_ip_limiter.allow(state.client_ip)
|
|
):
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "error",
|
|
"code": "rate_limited",
|
|
"message": "The channel needs a moment to settle before it can be focused again.",
|
|
}
|
|
)
|
|
return
|
|
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:
|
|
# Locked — see the same comment in _reward_summon above.
|
|
user = await db.scalar(select(User).where(User.id == state.user_id).with_for_update())
|
|
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
|
|
if not (
|
|
judgment_limiter.allow(str(state.user_id))
|
|
and judgment_ip_limiter.allow(state.client_ip)
|
|
):
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "error",
|
|
"code": "rate_limited",
|
|
"message": "The veil needs a moment before it can render another verdict.",
|
|
}
|
|
)
|
|
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:
|
|
# Locked — see the same comment in _reward_summon above; this
|
|
# path writes essence too, so it's exposed to the same race.
|
|
user = await db.scalar(select(User).where(User.id == state.user_id).with_for_update())
|
|
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
|