`passage_start` rewinds the rite to `listen` at any time, and PassagePanel offers exactly that button after a resisted release. Every replayed layer re-credited its essence, so one summon funded an endless loop — the 20/60s limiter caps the rate, never the total. Measured at ~110 essence/minute, indefinitely. Each layer now pays the first time it opens for a presence and never again, cleared only by a genuine summon. Re-walking still reveals; it just doesn't mint. The frame reports what was ACTUALLY credited, so the UI's running total can't drift from the ledger. Three further fixes in the same handlers: - judgment -> passage double-paid a crossing. The passage -> cross_over direction was already guarded; the reverse ran free, favor included. - both handlers read `state.entity["id"]` AFTER their DB round-trip. The HTTP telemetry path drives the same SeanceState and can summon concurrently, so a crossing could mark the presence that just arrived. Pinned before the awaits. tests/test_ws_passage.py is new, and covers the gap that let all of this hide: test_passage.py tests the pure module, and nothing exercised these handlers over a real connection. Efficacy proven by reverting the fix — the replay test then reports "minted 280 extra essence". Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
1344 lines
56 KiB
Python
1344 lines
56 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
|
|
{"type": "passage_start"} → passage_state (the rite begins)
|
|
{"type": "passage_layer"} → passage_result (one beat of the
|
|
layered crossing rite: listen → name → unbind → open → release)
|
|
{"type": "scry", "image": "<base64 jpeg>"} → {"type": "utterance", kind: "scry"}
|
|
The seeker's camera, shown to the vision model so the entity can speak
|
|
about the real room. The frame is never stored or logged — only the
|
|
resulting utterance is, like any other spirit speech.
|
|
|
|
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, passage
|
|
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.celestial import veil_thinness
|
|
from app.entropy import contribution_bits, veil_float
|
|
from app.geomagnetic import geomagnetic_cache
|
|
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", "camera"}
|
|
|
|
# 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)
|
|
# Scrying sends a real image to a vision model — by far the heaviest
|
|
# request this app makes of a CPU-only Ollama box, so it gets the
|
|
# tightest budget of any channel.
|
|
scry_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)
|
|
ritual_ip_limiter = RateLimiter(max_requests=12, window_seconds=60)
|
|
judgment_ip_limiter = RateLimiter(max_requests=20, window_seconds=60)
|
|
# The Passage (Workstream L): each `passage_layer` frame is its own essence
|
|
# credit, so it needs the same bounding as ritual/judgment. The budget is
|
|
# larger because one rite is five frames and a collapse legitimately makes a
|
|
# seeker replay earlier layers — four full rites a minute is already far
|
|
# past human pace, but never trips on honest play.
|
|
passage_limiter = RateLimiter(max_requests=20, window_seconds=60)
|
|
passage_ip_limiter = RateLimiter(max_requests=40, window_seconds=60)
|
|
scry_ip_limiter = RateLimiter(max_requests=8, window_seconds=60)
|
|
|
|
# Probability that a channel's familiar presence answers again rather than
|
|
# something new manifesting. High enough that the Codex stays collectable
|
|
# and spirits are genuinely re-contactable; low enough that calling into a
|
|
# known channel is never a guarantee.
|
|
RETURN_CHANCE = 0.72
|
|
|
|
# How much a fully-thin veil erodes the familiar presence's claim on a
|
|
# channel. At 0.45, a full moon at true solar midnight drops the return
|
|
# chance from 72% to ~40% — a real, felt difference on the spookiest night
|
|
# of the month, without ever making a known spirit unreachable.
|
|
VEIL_THINNESS_PULL = 0.45
|
|
|
|
|
|
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
|
|
# Workstream L (passage-doctrine spec): where the layered crossing rite
|
|
# stands for the *current* entity. `passage_layer` is the next beat to
|
|
# attempt, `passage_lied` remembers whether the layer just completed was
|
|
# a lie (so the next one can flag it in hindsight), and `passage_crossed`
|
|
# latches once the spirit has actually crossed so the rite — and the
|
|
# cross_over verdict — can't be replayed for a second payout. All three
|
|
# reset on a fresh summon, like the ritual fields above.
|
|
passage_layer: str = passage.FIRST_LAYER
|
|
passage_lied: bool = False
|
|
passage_crossed: bool = False
|
|
# Which beats have ALREADY paid out for the current entity. Without
|
|
# this, essence is unbounded: `passage_start` rewinds the rite to
|
|
# `listen` at any time (and the UI offers exactly that button after a
|
|
# resisted release), so a client could walk listen/name/unbind/open for
|
|
# +28 essence, restart, and repeat forever off a single summon — the
|
|
# limiter caps the rate, not the total. A volatility collapse rewinds it
|
|
# the same way. So each layer pays the FIRST time it opens for a given
|
|
# presence and never again; re-walking it still reveals, still costs
|
|
# nothing already banked, but mints no new essence. Cleared only on a
|
|
# fresh summon, never by `passage_start` — that is the whole point.
|
|
passage_paid: set[str] = field(default_factory=set)
|
|
# 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)
|
|
# Latest physical-entropy contribution harvested from the seeker's room
|
|
# (microphone noise floor / RF noise between stations) — see
|
|
# app/entropy.py. Untrusted by construction: it is only ever mixed with
|
|
# fresh server secrets, never used as a seed on its own, so a client
|
|
# sending a chosen or replayed value cannot steer any outcome.
|
|
entropy: str | None = None
|
|
# Seeker's longitude, if they granted location. Only the longitude is
|
|
# kept — it is all that solar midnight needs, and storing a full
|
|
# coordinate would be retaining precise location data we have no use
|
|
# for. Never persisted; lives and dies with the connection.
|
|
longitude: float | None = None
|
|
# Latest NOAA planetary K-index reading, refreshed on summon. Cached on
|
|
# the state so the manifest path can report it without another lookup.
|
|
geomagnetic: dict | None = None
|
|
# Serialises summoning for this session. Two independent coroutines can
|
|
# drive the same SeanceState: the browser's own WS message loop, and the
|
|
# HTTP telemetry ingestion path (device_anomaly.process_device_reading_
|
|
# for_summon -> _handle_anomaly), which is the entire point of the ESP32
|
|
# integration. Without this, both can observe `state.entity is None`,
|
|
# both await a summon that includes a multi-second LLM mint, and the
|
|
# session ends up with two entities minted, two essence credits, two
|
|
# item rolls, and whichever finishes last clobbering state.entity.
|
|
summon_lock: asyncio.Lock = field(default_factory=asyncio.Lock)
|
|
|
|
|
|
# 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]:
|
|
"""Open a channel and see what answers.
|
|
|
|
The signature identifies a *channel*, not a spirit. Whether the entity
|
|
previously reached on this channel answers again, or something else
|
|
manifests instead, is a genuine draw against physical entropy harvested
|
|
from the seeker's room (app/entropy.py) — not a deterministic lookup.
|
|
Before this, an identical anomaly pattern always produced an identical
|
|
spirit, which made contact feel like a database query rather than
|
|
channeling.
|
|
|
|
The Codex mechanic survives: a familiar presence is the *likely*
|
|
outcome on a known channel (`RETURN_CHANCE`), so spirits remain
|
|
collectable and re-contactable, but never guaranteed. Sometimes you
|
|
call and something else picks up.
|
|
|
|
An at-peace entity (a spirit correctly helped to cross over) is
|
|
excluded from the match entirely — it stays in the Codex forever but
|
|
can't be re-contacted. `Entity.signature` is unique, so a newly minted
|
|
entity can't reuse a string a retired row still holds; it gets a salted
|
|
variant of the same base signature.
|
|
"""
|
|
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:
|
|
known = await db.scalar(
|
|
select(Entity).where(Entity.signature == signature, Entity.at_peace.is_(False))
|
|
)
|
|
# The draw that makes contact feel like contact: even on a
|
|
# channel with a familiar presence, something else can
|
|
# answer. Domain-separated from the traits roll so the two
|
|
# are uncorrelated.
|
|
# A thinner veil (full moon, true solar midnight) makes it
|
|
# likelier that something *other* than the familiar
|
|
# presence pushes through — more traffic gets across when
|
|
# the barrier is weaker, which is the whole folkloric
|
|
# premise. Real astronomy, computed from the seeker's own
|
|
# longitude: see app/celestial.py.
|
|
sky = veil_thinness(state.longitude)
|
|
thinness = sky["thinness"]
|
|
if state.geomagnetic is not None:
|
|
# A real storm counts for a third of the reading — enough
|
|
# that a Kp-7 night is noticeably stranger, not so much
|
|
# that quiet space weather flattens the moon's effect.
|
|
thinness = thinness * 0.67 + state.geomagnetic["disturbance"] * 0.33
|
|
return_chance = RETURN_CHANCE * (1.0 - thinness * VEIL_THINNESS_PULL)
|
|
answers = (
|
|
known is not None
|
|
and veil_float(state.entropy, "answers") < return_chance
|
|
)
|
|
entity = known if answers else None
|
|
is_new = entity is None
|
|
|
|
if is_new:
|
|
# A new presence on an occupied channel needs its own
|
|
# signature (the column is unique) — same salting the
|
|
# at-peace case already required.
|
|
mint_signature = signature
|
|
taken = await db.scalar(select(Entity).where(Entity.signature == signature))
|
|
if taken is not None:
|
|
mint_signature = f"{signature}:{uuid.uuid4().hex[:8]}"
|
|
|
|
profile = await spirit_service.mint_profile(
|
|
mint_signature, channel, state.anomalies, state.language,
|
|
entropy=state.entropy,
|
|
# The moon shapes who answers — see
|
|
# entities.rarity_weights_for_moon / moon_trait_bias.
|
|
sky=sky,
|
|
)
|
|
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:
|
|
# Held across the whole mint so a concurrent caller waits rather than
|
|
# starting a second summon. See SeanceState.summon_lock.
|
|
async with state.summon_lock:
|
|
await _summon_locked(state)
|
|
|
|
|
|
async def _summon_locked(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
|
|
|
|
# Real measured geomagnetic activity (NOAA SWPC). Served from cache and
|
|
# never blocking: a cold cache or an outage yields None and the séance
|
|
# proceeds on astronomy alone.
|
|
geo = await geomagnetic_cache.reading()
|
|
if geo is not None:
|
|
state.geomagnetic = geo
|
|
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "status",
|
|
"state": "summoning",
|
|
"geomagnetic": state.geomagnetic,
|
|
# How much physical noise from the room actually fed this
|
|
# draw. Display only — a client can lie about it freely, and
|
|
# it is never used to weight or gate anything.
|
|
"entropy_bits": contribution_bits(state.entropy),
|
|
# Real astronomy, computed not fetched (app/celestial.py) — the
|
|
# seeker can see *why* tonight is different.
|
|
"sky": veil_thinness(state.longitude),
|
|
}
|
|
)
|
|
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
|
|
_reset_passage(state, new_entity=True)
|
|
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
|
|
|
|
|
|
# Chance that a shift in the room pulls unprompted speech through. This is
|
|
# not a reply to anything — see SpiritService.manifest(): the prompt
|
|
# contains no seeker input at all, and the token-sampling seed comes from
|
|
# physical noise measured in that room. Kept well under half so silence
|
|
# stays the norm and being spoken to unbidden stays unnerving rather than
|
|
# chatty.
|
|
MANIFEST_CHANCE_ON_ANOMALY = 0.28
|
|
|
|
|
|
def _room_readings(state: SeanceState, anomaly: dict) -> dict:
|
|
"""The measured state of the room, rendered plainly for the entity.
|
|
|
|
Values are reported as measurements, never as interpretations — the
|
|
model should react to "magnitude 31.4 dB over floor", not to
|
|
"terrifying paranormal spike". Putting the conclusion in the prompt
|
|
would mean the horror came from us instead of from the entity.
|
|
"""
|
|
readings: dict[str, str] = {}
|
|
source = anomaly.get("source")
|
|
if source:
|
|
readings["channel"] = str(source)
|
|
freq = anomaly.get("frequency")
|
|
if isinstance(freq, (int, float)):
|
|
readings["frequency"] = f"{freq:.2f}"
|
|
mag = anomaly.get("magnitude")
|
|
if isinstance(mag, (int, float)):
|
|
readings["deviation above the floor"] = f"{mag:.1f}"
|
|
readings["disturbances so far"] = str(len(state.anomalies))
|
|
if state.geomagnetic is not None:
|
|
readings["earth's magnetic field"] = (
|
|
f"Kp {state.geomagnetic['kp']:.2f} — {state.geomagnetic['label']}"
|
|
)
|
|
sky = veil_thinness(state.longitude)
|
|
readings["moon"] = f"{sky['moon_name']} ({sky['moon_illumination']:.0%} lit)"
|
|
if sky["witching_proximity"] is not None:
|
|
readings["nearness to the dead of night"] = f"{sky['witching_proximity']:.0%}"
|
|
if state.entropy:
|
|
readings["noise gathered from the room"] = f"{contribution_bits(state.entropy)} bits"
|
|
return readings
|
|
|
|
|
|
async def _maybe_manifest(state: SeanceState, anomaly: dict) -> None:
|
|
"""Let the entity speak unbidden, if the room pulls it through."""
|
|
if state.entity is None:
|
|
return
|
|
if state.tell_rng.random() >= MANIFEST_CHANCE_ON_ANOMALY:
|
|
return
|
|
try:
|
|
text = await spirit_service.manifest(
|
|
state.entity,
|
|
_room_readings(state, anomaly),
|
|
state.language,
|
|
entropy=state.entropy,
|
|
)
|
|
except Exception:
|
|
# Deliberately broad. Unprompted speech is a bonus, never
|
|
# load-bearing — a busy queue, an unreachable Ollama, or a
|
|
# malformed response must leave the séance quiet rather than
|
|
# surfacing an error for something the seeker never asked for.
|
|
return
|
|
if text:
|
|
await _speak(state, "manifest", text)
|
|
|
|
|
|
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)
|
|
await _maybe_manifest(state, anomaly)
|
|
|
|
|
|
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
|
|
|
|
# The Passage already crossed this spirit and already paid for it — a
|
|
# `cross_over` verdict on top would credit CROSS_OVER_ESSENCE a second
|
|
# time for the same crossing. One spirit, one passage, one payment.
|
|
if verdict == "cross_over" and state.passage_crossed:
|
|
return
|
|
|
|
traits = state.entity.get("traits", {})
|
|
# Pinned before any await, for the same reason as in the Passage handler
|
|
# below: the device-telemetry path can re-summon on this SeanceState
|
|
# mid-handler, and re-reading `state.entity` after the DB round-trip would
|
|
# mark whichever entity arrived last as at_peace instead of the judged one.
|
|
entity_id = state.entity["id"]
|
|
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(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()
|
|
|
|
if outcome.consequence == "crossed_over":
|
|
# Latch the same flag the Passage sets. The guard above already stops
|
|
# passage → cross_over double-paying; without this the REVERSE ran
|
|
# free: judge cross_over (+essence, +favor, entity at_peace), then walk
|
|
# the Passage to `release` on the same spirit and be paid the crossing
|
|
# a second time, favor included. One spirit, one crossing, whichever
|
|
# road got there first.
|
|
state.passage_crossed = True
|
|
|
|
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})
|
|
|
|
|
|
# --- Workstream L: the Passage (passage-doctrine spec) ---------------------
|
|
#
|
|
# Five beats — listen / name / unbind / open / release — each revealing one
|
|
# real trait, paying a little essence, and able to twist. All of the actual
|
|
# logic is in `app/passage.py`, which is pure; this handler owns only the
|
|
# things a pure module can't: the limiter, the two entropy draws, the
|
|
# essence credit (via the same `credit_essence` the rest of the economy
|
|
# uses), the `at_peace` write, and the frames.
|
|
|
|
|
|
def _reset_passage(state: SeanceState, *, new_entity: bool = False) -> None:
|
|
state.passage_layer = passage.FIRST_LAYER
|
|
state.passage_lied = False
|
|
state.passage_crossed = False
|
|
if new_entity:
|
|
# Only a genuinely new presence re-opens the purse. A `passage_start`
|
|
# rewind must not, or the rite becomes an essence faucet (see
|
|
# `passage_paid` on SeanceState).
|
|
state.passage_paid = set()
|
|
|
|
|
|
def _passage_frame(
|
|
state: SeanceState, outcome: passage.PassageOutcome, essence: int
|
|
) -> dict:
|
|
# `essence` is what was ACTUALLY credited, which is the layer's value only
|
|
# the first time it opens for this presence (see `passage_paid`). The
|
|
# client sums this field into its running total, so sending
|
|
# `outcome.essence` on a replay would show the seeker essence they did not
|
|
# receive.
|
|
reveal = outcome.reveal
|
|
return {
|
|
"type": "passage_result",
|
|
"layer": outcome.layer,
|
|
"result": outcome.result,
|
|
# The client renders the words from (trait, band) via i18n, so no
|
|
# untranslated prose ever crosses the wire.
|
|
"reveal": (
|
|
{"trait": reveal.trait, "value": reveal.value, "band": reveal.band}
|
|
if reveal is not None
|
|
else None
|
|
),
|
|
"essence": essence,
|
|
"at_peace": outcome.at_peace,
|
|
"next_layer": outcome.next_layer,
|
|
# Deliberately NOT `outcome.lied` — that's this layer's own lie, and
|
|
# the seeker must not learn of it until the next beat. This field is
|
|
# the *previous* layer's lie, surfacing in hindsight.
|
|
"unreliable_layer": outcome.unreliable_layer,
|
|
"sealed": state.passage_crossed,
|
|
}
|
|
|
|
|
|
async def _handle_passage_start(state: SeanceState) -> None:
|
|
if state.entity is None:
|
|
return # no presence to pass — the frontend already gates the button
|
|
if state.passage_crossed:
|
|
return # already at peace; there is nothing left to walk
|
|
if not (
|
|
passage_limiter.allow(str(state.user_id))
|
|
and passage_ip_limiter.allow(state.client_ip)
|
|
):
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "error",
|
|
"code": "rate_limited",
|
|
"message": "The way is still closing behind you. Give it a moment.",
|
|
}
|
|
)
|
|
return
|
|
_reset_passage(state)
|
|
await state.send_queue.put(
|
|
{"type": "passage_state", "layer": state.passage_layer, "sealed": False}
|
|
)
|
|
|
|
|
|
async def _handle_passage_layer(state: SeanceState) -> None:
|
|
if state.entity is None or state.passage_crossed:
|
|
return
|
|
if not (
|
|
passage_limiter.allow(str(state.user_id))
|
|
and passage_ip_limiter.allow(state.client_ip)
|
|
):
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "error",
|
|
"code": "rate_limited",
|
|
"message": "The way is still closing behind you. Give it a moment.",
|
|
}
|
|
)
|
|
return
|
|
|
|
traits = state.entity.get("traits", {})
|
|
# Pinned HERE, before any await. `state.entity` can be replaced underneath
|
|
# this coroutine mid-handler: the HTTP device-telemetry path drives the
|
|
# same SeanceState and can summon (see `summon_lock`). Reading the id
|
|
# after the DB round-trip below would mark the WRONG entity at peace —
|
|
# the one that just arrived, not the one that actually crossed.
|
|
entity_id = state.entity["id"]
|
|
layer = state.passage_layer
|
|
# Both twists draw from the room's real physical noise, with distinct
|
|
# contexts so a volatility collapse and a deceptive reveal can never be
|
|
# correlated with each other (or with any other draw in the app).
|
|
draw = passage.PassageDraw(
|
|
collapse=veil_float(state.entropy, f"passage:collapse:{layer}"),
|
|
deceit=veil_float(state.entropy, f"passage:deceit:{layer}"),
|
|
)
|
|
outcome = passage.resolve_layer(
|
|
traits, layer, draw, previous_lied=state.passage_lied
|
|
)
|
|
|
|
state.passage_lied = outcome.lied
|
|
if outcome.result == "collapsed":
|
|
# The rite restarts from `listen` — but the essence already earned
|
|
# stays credited. Taking it back would punish the seeker for the
|
|
# spirit's own instability.
|
|
state.passage_lied = False
|
|
state.passage_layer = outcome.next_layer or layer
|
|
if outcome.at_peace:
|
|
state.passage_crossed = True
|
|
|
|
# A layer pays once per presence. Replaying it — via `passage_start`, via
|
|
# a collapse rewind, or via a hand-rolled client spamming the frame —
|
|
# reveals the same thing again for free rather than minting essence again.
|
|
award = outcome.essence if layer not in state.passage_paid else 0
|
|
if outcome.result in ("opened", "crossed"):
|
|
state.passage_paid.add(layer)
|
|
|
|
item = None
|
|
if award or outcome.at_peace:
|
|
async with session_maker() as db:
|
|
# Locked — same read-modify-write race as every other essence
|
|
# credit in this file (see _reward_summon).
|
|
user = await db.scalar(
|
|
select(User).where(User.id == state.user_id).with_for_update()
|
|
)
|
|
if user is not None:
|
|
if award:
|
|
credit_essence(user, award)
|
|
if outcome.at_peace:
|
|
# Completing the Passage is the most compassionate act
|
|
# in the game, exactly as a correct cross_over verdict
|
|
# is — so it moves favor by the same amount.
|
|
user.favor = judgment.clamp_favor(
|
|
user.favor + judgment.FAVOR_CORRECT_CROSS_OVER
|
|
)
|
|
if outcome.at_peace:
|
|
entity_row = await db.get(Entity, uuid.UUID(entity_id))
|
|
if entity_row is not None:
|
|
entity_row.at_peace = True
|
|
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(_passage_frame(state, outcome, award))
|
|
if item is not None:
|
|
await state.send_queue.put({"type": "item_drop", "item": item})
|
|
|
|
|
|
# A 768px JPEG at quality 0.72 is well under 200KB, so ~350KB of base64 is a
|
|
# generous ceiling that still refuses anything pathological before it reaches
|
|
# the model.
|
|
MAX_SCRY_B64_CHARS = 350_000
|
|
|
|
|
|
async def _handle_scry(state: SeanceState, message: dict) -> None:
|
|
"""The entity speaks about what the seeker's camera actually shows.
|
|
|
|
Unlike every other channel, the payload here is a photograph of a real
|
|
room, so this handler is deliberately strict: no entity means nothing to
|
|
look through, the image is size-capped before it touches the queue, and
|
|
the frame is never persisted or logged anywhere — it is passed to the
|
|
model and dropped. Only the resulting utterance is recorded, exactly like
|
|
any other thing a spirit says.
|
|
"""
|
|
if state.entity is None:
|
|
return
|
|
image = message.get("image")
|
|
if not isinstance(image, str) or not image.strip():
|
|
return
|
|
if len(image) > MAX_SCRY_B64_CHARS:
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "error",
|
|
"code": "scry_too_large",
|
|
"message": "the lens showed too much at once — try again.",
|
|
}
|
|
)
|
|
return
|
|
if not (
|
|
scry_limiter.allow(str(state.user_id))
|
|
and scry_ip_limiter.allow(state.client_ip)
|
|
):
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "error",
|
|
"code": "rate_limited",
|
|
"message": "the eye tires. let it rest a moment before looking again.",
|
|
}
|
|
)
|
|
return
|
|
|
|
await state.send_queue.put({"type": "status", "state": "gathering"})
|
|
try:
|
|
text = await spirit_service.scry(
|
|
state.entity, image, state.language, entropy=state.entropy
|
|
)
|
|
except SpiritBusyError:
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "error",
|
|
"code": "veil_crowded",
|
|
"message": "too many eyes at once. try again shortly.",
|
|
}
|
|
)
|
|
return
|
|
except Exception:
|
|
# Deliberately broad, same reasoning as _maybe_manifest: a vision
|
|
# failure must leave the séance intact rather than surfacing a stack
|
|
# trace for something the seeker can't act on.
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "error",
|
|
"code": "scry_failed",
|
|
"message": "the lens clouded over. nothing came through.",
|
|
}
|
|
)
|
|
return
|
|
if text:
|
|
await _speak(state, "scry", text)
|
|
|
|
|
|
@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":
|
|
# The seeker's room contributes its physical noise to this
|
|
# draw. Stored raw and untrusted — app/entropy.py mixes it
|
|
# with fresh server secrets on every use, so a chosen or
|
|
# replayed value can add unpredictability but never steer
|
|
# the outcome.
|
|
contributed = message.get("entropy")
|
|
if isinstance(contributed, str) and contributed.strip():
|
|
state.entropy = contributed[:512]
|
|
# Longitude only, and only if the seeker granted location.
|
|
# Out-of-range values are dropped rather than clamped: a
|
|
# bogus longitude should mean "no location", not a
|
|
# confidently wrong solar midnight.
|
|
lon = message.get("longitude")
|
|
if isinstance(lon, (int, float)) and -180 <= lon <= 180:
|
|
state.longitude = float(lon)
|
|
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)
|
|
elif msg_type == "passage_start":
|
|
await _handle_passage_start(state)
|
|
elif msg_type == "passage_layer":
|
|
await _handle_passage_layer(state)
|
|
elif msg_type == "scry":
|
|
await _handle_scry(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
|