Adds backend/app/device_anomaly.py — per-(user_id, device_id, sensor_type) rolling-baseline anomaly detection for continuous numeric sensors (structurally modeled on telemetry.detect_wire_spike: min samples, an absolute floor, 3-sigma + relative threshold, with per-sensor-type floors since units vary wildly) plus a false->true state-transition detector for discrete/boolean sensors like presence. Adds a module-level active-session registry in app/ws.py (register_active_session/unregister_active_session/get_active_session) so hardware ingestion (a plain HTTP call, not a WS connection) can find a user's live SeanceState. process_device_reading_for_summon(user_id, device_id, sensor_type, value, unit) is the self-contained entry point Workstream G's ingestion handler will call into: classifies numeric vs. boolean, runs the reading through the right detector, and on a genuine anomaly pushes it into the active session via the existing _handle_anomaly path (source=sensor_type, frequency=stable per-sensor-type constant, magnitude=deviation-from- baseline or a fixed constant for boolean transitions) — reusing the full existing signature/mint/Codex pipeline, no new mint logic. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
523 lines
20 KiB
Python
523 lines
20 KiB
Python
"""The séance channel: one WebSocket per contact session.
|
|
|
|
Protocol (client → server):
|
|
{"type": "ping"} → {"type": "pong"}
|
|
{"type": "set_mode", "mode": MODE} → {"type": "mode", ...}
|
|
{"type": "language", "language": "en"|"es"}
|
|
{"type": "summon"} → {"type": "entity", ...}
|
|
{"type": "anomaly", "source": SRC, ...} → {"type": "utterance", ...}
|
|
{"type": "question", "text": "..."} → reply_start / reply_token* / reply_end
|
|
{"type": "passive", "enabled": bool} → ambient wire loop on/off
|
|
|
|
All server → client frames flow through a single sender task so concurrent
|
|
producers (ambient loop, reply streaming, TTS callbacks) never interleave on
|
|
the wire.
|
|
"""
|
|
|
|
import asyncio
|
|
import contextlib
|
|
import random
|
|
import time
|
|
import uuid
|
|
from dataclasses import dataclass, field
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
|
|
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
|
|
from sqlalchemy import select
|
|
|
|
from app.config import settings
|
|
from app.db import async_session_maker as _default_session_maker
|
|
from app.deps import SESSION_COOKIE_NAME
|
|
from app.entities import fallback_signature, signature_from_anomalies
|
|
from app.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.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
|
|
|
|
|
|
# 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(),
|
|
}
|
|
|
|
|
|
def _client_ip(websocket: WebSocket) -> str:
|
|
host = websocket.client.host if websocket.client else None
|
|
return resolve_client_ip(websocket.headers, host)
|
|
|
|
|
|
async def _authenticate(websocket: WebSocket) -> uuid.UUID | None:
|
|
raw_token = websocket.cookies.get(SESSION_COOKIE_NAME)
|
|
if raw_token is None:
|
|
return None
|
|
|
|
token_hash = hash_token(raw_token)
|
|
async with session_maker() as db:
|
|
result = await db.execute(
|
|
select(AuthSession).where(AuthSession.token_hash == token_hash)
|
|
)
|
|
session = result.scalar_one_or_none()
|
|
if session is None or session.expires_at < datetime.now(timezone.utc):
|
|
return None
|
|
return session.user_id
|
|
|
|
|
|
async def _sender(state: SeanceState, websocket: WebSocket) -> None:
|
|
"""The only task allowed to write to the socket."""
|
|
try:
|
|
while True:
|
|
message = await state.send_queue.get()
|
|
await websocket.send_json(message)
|
|
except asyncio.CancelledError:
|
|
pass
|
|
|
|
|
|
async def _record_event(
|
|
session_id: uuid.UUID,
|
|
kind: str,
|
|
text: str | None = None,
|
|
payload: dict | None = None,
|
|
) -> uuid.UUID:
|
|
async with session_maker() as db:
|
|
event = Event(session_id=session_id, kind=kind, text=text, payload=payload)
|
|
db.add(event)
|
|
await db.commit()
|
|
await db.refresh(event)
|
|
return event.id
|
|
|
|
|
|
async def _speak(
|
|
state: SeanceState, kind: str, text: str, instability: float = 0.0
|
|
) -> None:
|
|
"""Persist an utterance, push its text immediately, synthesize audio in
|
|
the background, and push the audio URL when the effects chain finishes."""
|
|
event_id = await _record_event(
|
|
state.session_id, "utterance", text=text, payload={"kind": kind}
|
|
)
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "utterance",
|
|
"id": str(event_id),
|
|
"kind": kind,
|
|
"text": text,
|
|
"entity": state.entity["name"] if state.entity else None,
|
|
}
|
|
)
|
|
|
|
async def _synth() -> None:
|
|
try:
|
|
voice_profile = state.entity.get("voice", {}) if state.entity else {}
|
|
voice = pick_voice(voice_profile.get("voice_id"), state.language)
|
|
wav = await synthesize_spirit_voice(
|
|
text, voice, voice_profile, instability=instability
|
|
)
|
|
filename = f"{event_id}.wav"
|
|
_audio_dir().joinpath(filename).write_bytes(wav)
|
|
async with session_maker() as db:
|
|
event = await db.get(Event, event_id)
|
|
if event is not None:
|
|
event.audio_path = filename
|
|
await db.commit()
|
|
await state.send_queue.put(
|
|
{"type": "audio", "id": str(event_id), "url": f"/audio/{filename}"}
|
|
)
|
|
except Exception:
|
|
# TTS is texture, not content — the words already reached the
|
|
# seeker. Never let a synth failure kill the séance.
|
|
pass
|
|
|
|
asyncio.create_task(_synth())
|
|
|
|
|
|
_ROMAN = ["II", "III", "IV", "V", "VI", "VII", "VIII", "IX"]
|
|
|
|
|
|
async def _unique_entity_name(db, base_name: str) -> str:
|
|
name = base_name
|
|
i = 0
|
|
while await db.scalar(select(Entity).where(Entity.name == name)) is not None:
|
|
name = f"{base_name} {_ROMAN[i]}" if i < len(_ROMAN) else f"{base_name} {i + 2}"
|
|
i += 1
|
|
return name
|
|
|
|
|
|
async def _summon(state: SeanceState, channel: str) -> tuple[Entity, bool]:
|
|
"""Match this session's signature against the Codex, or mint a new entity."""
|
|
signature = signature_from_anomalies(state.anomalies) or fallback_signature(
|
|
str(state.session_id)
|
|
)
|
|
|
|
async with session_maker() as db:
|
|
entity = await db.scalar(select(Entity).where(Entity.signature == signature))
|
|
is_new = entity is None
|
|
|
|
if is_new:
|
|
profile = await spirit_service.mint_profile(
|
|
signature, channel, state.anomalies, state.language
|
|
)
|
|
entity = Entity(
|
|
name=await _unique_entity_name(db, profile["name"]),
|
|
epithet=profile["epithet"],
|
|
persona=profile["persona"],
|
|
rarity_tier=profile["rarity"],
|
|
signature=signature,
|
|
voice_profile=profile["voice"],
|
|
visual_profile=profile["visual"],
|
|
sample_quotes=profile["quotes"],
|
|
discovered_by=state.user_id,
|
|
contact_count=1,
|
|
)
|
|
db.add(entity)
|
|
else:
|
|
entity.contact_count += 1
|
|
|
|
await db.flush()
|
|
session = await db.get(ContactSession, state.session_id)
|
|
if session is not None:
|
|
session.entity_id = entity.id
|
|
db.add(
|
|
EntitySighting(
|
|
entity_id=entity.id, session_id=state.session_id, user_id=state.user_id
|
|
)
|
|
)
|
|
await db.commit()
|
|
await db.refresh(entity)
|
|
return entity, is_new
|
|
|
|
|
|
async def _handle_summon(state: SeanceState) -> None:
|
|
# Short-circuits: an account already over its own cap never gets far
|
|
# enough to spend from the IP budget too.
|
|
if not (
|
|
summon_limiter.allow(str(state.user_id))
|
|
and summon_ip_limiter.allow(state.client_ip)
|
|
):
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "error",
|
|
"code": "rate_limited",
|
|
"message": "The veil is crowded. The spirits need a moment before another summoning.",
|
|
}
|
|
)
|
|
return
|
|
|
|
await state.send_queue.put({"type": "status", "state": "summoning"})
|
|
entity, is_new = await _summon(state, state.mode if state.mode != "unknown" else "ouija")
|
|
state.entity = serialize_entity(entity)
|
|
await state.send_queue.put(
|
|
{"type": "entity", "entity": state.entity, "is_new": is_new}
|
|
)
|
|
greeting = random.choice(state.entity["quotes"]) if state.entity["quotes"] else "I am here."
|
|
await _speak(state, "greeting", greeting)
|
|
|
|
|
|
async def _handle_anomaly(state: SeanceState, message: dict) -> None:
|
|
anomaly = {
|
|
"source": str(message.get("source", "unknown"))[:16],
|
|
"frequency": message.get("frequency"),
|
|
"magnitude": message.get("magnitude"),
|
|
}
|
|
state.anomalies.append(anomaly)
|
|
state.anomalies = state.anomalies[-64:]
|
|
await _record_event(state.session_id, "anomaly", payload=anomaly)
|
|
await state.send_queue.put({"type": "anomaly_ack", "count": len(state.anomalies)})
|
|
|
|
if state.entity is None:
|
|
if signature_from_anomalies(state.anomalies) is not None:
|
|
# The signal has enough structure — something announces itself.
|
|
await _handle_summon(state)
|
|
else:
|
|
await state.send_queue.put({"type": "status", "state": "attuning"})
|
|
return
|
|
|
|
if not (
|
|
fragment_limiter.allow(str(state.user_id))
|
|
and fragment_ip_limiter.allow(state.client_ip)
|
|
):
|
|
return # anomalies during a crowded veil just pass unheard
|
|
|
|
try:
|
|
fragment = await spirit_service.fragment(
|
|
anomaly["source"], anomaly, state.language
|
|
)
|
|
except SpiritBusyError:
|
|
return
|
|
await _speak(state, "fragment", fragment)
|
|
|
|
|
|
async def _handle_question(state: SeanceState, text: str) -> None:
|
|
if not (
|
|
question_limiter.allow(str(state.user_id))
|
|
and question_ip_limiter.allow(state.client_ip)
|
|
):
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "error",
|
|
"code": "rate_limited",
|
|
"message": "The spirit is spent. Give it a moment to gather itself.",
|
|
}
|
|
)
|
|
return
|
|
|
|
if state.entity is None:
|
|
await _handle_summon(state)
|
|
assert state.entity is not None
|
|
|
|
text = text.strip()[:500]
|
|
await _record_event(state.session_id, "question", text=text)
|
|
await state.send_queue.put({"type": "status", "state": "gathering"})
|
|
|
|
last_magnitude = state.anomalies[-1].get("magnitude") if state.anomalies else 50
|
|
if last_magnitude is None:
|
|
last_magnitude = 50
|
|
stability = compute_stability(state.entity["rarity"], float(last_magnitude))
|
|
await state.send_queue.put({"type": "reply_start", "stability": stability})
|
|
|
|
reply_parts: list[str] = []
|
|
try:
|
|
async for token in spirit_service.chat_stream(
|
|
state.entity, text, state.history, state.language
|
|
):
|
|
reply_parts.append(token)
|
|
await state.send_queue.put({"type": "reply_token", "token": token})
|
|
except SpiritBusyError:
|
|
await state.send_queue.put(
|
|
{
|
|
"type": "error",
|
|
"code": "veil_crowded",
|
|
"message": "Too many seekers press against the veil. The spirit withdraws.",
|
|
}
|
|
)
|
|
await state.send_queue.put({"type": "reply_end", "text": ""})
|
|
return
|
|
|
|
reply = "".join(reply_parts).strip()
|
|
reply_id = await _record_event(state.session_id, "reply", text=reply)
|
|
state.history.append({"role": "user", "text": text})
|
|
state.history.append({"role": "spirit", "text": reply})
|
|
state.history = state.history[-8:]
|
|
await state.send_queue.put({"type": "reply_end", "id": str(reply_id), "text": reply})
|
|
if reply:
|
|
await _speak(state, "reply", reply, instability=1 - stability)
|
|
|
|
|
|
async def _ambient_loop(state: SeanceState) -> None:
|
|
"""The Wire Ghost's pulse: telemetry every few seconds. Spikes in the
|
|
jitter baseline become first-class anomaly events (driving fragments and
|
|
the client's data views); quieter ticks may earn an ambient whisper."""
|
|
try:
|
|
while True:
|
|
await asyncio.sleep(random.uniform(3, 5))
|
|
try:
|
|
sample = await sample_network(period_s=1.0)
|
|
except Exception:
|
|
continue
|
|
await state.send_queue.put({"type": "telemetry", **sample.as_dict()})
|
|
|
|
jitter = sample.jitter_bytes_per_s
|
|
ratio = detect_wire_spike(state.wire_jitter_history, jitter)
|
|
state.wire_jitter_history = (state.wire_jitter_history + [jitter])[-60:]
|
|
|
|
now = time.monotonic()
|
|
if ratio is not None and now - state.last_wire_anomaly_at > 20:
|
|
state.last_wire_anomaly_at = now
|
|
anomaly = {
|
|
"source": "wire",
|
|
"frequency": round(jitter, 1),
|
|
"magnitude": round(ratio * 10, 1),
|
|
}
|
|
# Tell the client first so its matrix/graph views light up,
|
|
# then run the anomaly through the usual séance pipeline.
|
|
await state.send_queue.put({"type": "anomaly", **anomaly})
|
|
await _handle_anomaly(state, anomaly)
|
|
continue
|
|
|
|
if not spirit_service.ambient_ready():
|
|
continue
|
|
whisper = await spirit_service.wire_whisper(sample.as_dict(), state.language)
|
|
await _speak(state, "ambient", whisper)
|
|
except asyncio.CancelledError:
|
|
pass
|
|
|
|
|
|
async def _handle_passive(state: SeanceState, enabled: bool) -> None:
|
|
if enabled and (state.ambient_task is None or state.ambient_task.done()):
|
|
state.ambient_task = asyncio.create_task(_ambient_loop(state))
|
|
await state.send_queue.put({"type": "passive", "enabled": True})
|
|
elif not enabled and state.ambient_task is not None:
|
|
state.ambient_task.cancel()
|
|
state.ambient_task = None
|
|
await state.send_queue.put({"type": "passive", "enabled": False})
|
|
|
|
|
|
@router.websocket("/ws/session")
|
|
async def session_socket(websocket: WebSocket) -> None:
|
|
user_id = await _authenticate(websocket)
|
|
if user_id is None:
|
|
await websocket.close(code=4401)
|
|
return
|
|
|
|
await websocket.accept()
|
|
|
|
async with session_maker() as db:
|
|
contact_session = ContactSession(user_id=user_id)
|
|
db.add(contact_session)
|
|
await db.commit()
|
|
await db.refresh(contact_session)
|
|
|
|
state = SeanceState(
|
|
user_id=user_id,
|
|
session_id=contact_session.id,
|
|
client_ip=_client_ip(websocket),
|
|
)
|
|
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")))
|
|
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
|