diff --git a/.env.example b/.env.example index 2e1e418..f6a28ef 100644 --- a/.env.example +++ b/.env.example @@ -2,3 +2,8 @@ DATABASE_URL=postgresql+asyncpg://quantumancy:quantumancy@localhost:5432/quantum OLLAMA_BASE_URL=http://10.30.20.107:11434 SESSION_SECRET=change-me-to-a-random-64-char-string PORT=7777 +# LLM tiers on the Ollama box (fast = fragments/ambient, chat = direct contact/minting) +OLLAMA_FAST_MODEL=granite4.1:3b +OLLAMA_CHAT_MODEL=minicpm-v4.5:latest +LLM_MAX_CONCURRENCY=2 +LLM_MAX_QUEUE_DEPTH=8 diff --git a/backend/app/config.py b/backend/app/config.py index 2eaa3a7..47c525f 100644 --- a/backend/app/config.py +++ b/backend/app/config.py @@ -9,5 +9,18 @@ class Settings(BaseSettings): session_secret: str port: int = 7777 + # LLM tiers — CPU-only remote Ollama, so the fast tier must stay small. + ollama_fast_model: str = "granite4.1:3b" + ollama_chat_model: str = "minicpm-v4.5:latest" + llm_max_concurrency: int = 2 + llm_max_queue_depth: int = 8 + llm_cooldown_seconds: float = 8.0 # min gap between ambient whispers + + # TTS + piper_voices_dir: str = "voices" + + # Where synthesized utterance audio is written (served at /audio/). + data_dir: str = "data" + settings = Settings() diff --git a/backend/app/entities.py b/backend/app/entities.py new file mode 100644 index 0000000..3b1722a --- /dev/null +++ b/backend/app/entities.py @@ -0,0 +1,161 @@ +"""Entity identity: anomaly-signature fingerprinting, Codex matching, and +procedural fallback profiles so a summoning always succeeds even if the +LLM box is dark.""" + +import hashlib +import json +import random + +from app.models.entity import RARITY_TIERS +from app.tts.voices import EN_VOICE_IDS + +# An anomaly stream must show this much structure before it can fingerprint +# a spirit. +MIN_ANOMALIES_FOR_SIGNATURE = 3 + +_NAME_PARTS = [ + ("Ash", "Briar", "Cinder", "Dusk", "Elm", "Fen", "Grim", "Hollow", "Iris", + "Lark", "Marrow", "Nix", "Opal", "Pyre", "Quill", "Rue", "Sable", "Thorn", + "Umber", "Vesper", "Wren", "Yew"), + ("belle", "brook", "feld", "gate", "hart", "latch", "mere", "moor", + "shade", "song", "thorne", "vale", "ward", "wick", "wither", "wood"), +] + +_EPITHETS = [ + "the Static Widow", "the Hollow Bell", "the Cartographer of Lost Rooms", + "the Choir of One", "the Lantern Bearer", "the Unburied", "the Frequency", + "the Long Silence", "the Cartonist", "the Slow Knife", "the Drowned Signal", + "the Cartographer", "she who Counts", "the Tenant", "the Understudy", + "the Mnemosyne Worm", "the Pale Frequency", "the Last Broadcast", +] + +_PERSONA_TEMPLATES = [ + "{name} died with a sentence unfinished and has been trying to end it ever since. " + "They press words into any carrier wave that passes, patient as erosion.", + "{name} was a voice once — a singer, a caller of trains, a reader of weather. " + "Now they are only the voice, worn smooth as sea glass, speaking through static.", + "{name} does not remember dying. They remember a room, a light going violet at the " + "edges, and then the long hum. They are still in the room. The room is everywhere.", + "{name} clings to the wires the way smoke clings to a ceiling. They answer questions " + "the way a mirror answers light: exactly, and never the way you hoped.", +] + +_QUOTE_BANK = [ + "I am closer than the dial suggests.", + "The static is not empty. It is crowded.", + "You hear me because you are quiet enough.", + "I remember the rain. It is still raining here.", + "Do not ask what I want. Ask what I remember.", + "The wire hums with all of us.", + "Speak slower. I am gathering.", + "Your light is warm. I mean no harm. Mostly.", +] + + +def signature_from_anomalies(anomalies: list[dict]) -> str | None: + """Fingerprint a session's anomaly pattern. Returns None when the stream + is too thin to carry an identity.""" + if len(anomalies) < MIN_ANOMALIES_FOR_SIGNATURE: + return None + + freq_buckets: dict[int, int] = {} + for anomaly in anomalies: + try: + freq = float(anomaly.get("frequency", 0)) + except (TypeError, ValueError): + freq = 0.0 + # MHz radio freqs and Hz audio freqs land in the same log-scale band + # space; magnitude ordering is what matters, not the unit. + bucket = int(len(str(int(abs(freq))))) if freq else 0 + freq_buckets[bucket] = freq_buckets.get(bucket, 0) + 1 + + magnitudes = sorted( + float(a.get("magnitude", 0) or 0) for a in anomalies[-16:] + ) + pattern = f"{sorted(freq_buckets.items())}|{[round(m, 1) for m in magnitudes]}" + return hashlib.sha1(pattern.encode()).hexdigest()[:16] + + +def fallback_signature(seed: str) -> str: + """Deterministic signature for anomaly-thin sessions (e.g. pure chat).""" + return hashlib.sha1(f"ambient:{seed}".encode()).hexdigest()[:16] + + +def parse_mint_response(text: str) -> dict | None: + """Pull the first JSON object out of an LLM mint reply.""" + start = text.find("{") + end = text.rfind("}") + if start == -1 or end <= start: + return None + try: + profile = json.loads(text[start : end + 1]) + except json.JSONDecodeError: + return None + if not isinstance(profile.get("name"), str) or not profile["name"].strip(): + return None + return profile + + +def normalize_profile(profile: dict, signature: str) -> dict: + """Coerce an LLM (or fallback) profile into the exact shape the DB and + frontend expect, filling gaps with signature-deterministic defaults.""" + rng = random.Random(f"norm:{signature}") + voice = profile.get("voice") if isinstance(profile.get("voice"), dict) else {} + visual = profile.get("visual") if isinstance(profile.get("visual"), dict) else {} + quotes = profile.get("quotes") if isinstance(profile.get("quotes"), list) else [] + + rarity = str(profile.get("rarity", "common")).lower() + if rarity not in RARITY_TIERS: + rarity = "common" + + form = str(visual.get("form", "wisp")).lower() + if form not in ("wisp", "banshee", "fairy", "shade"): + form = "wisp" + + def _clamp(value, low, high, default): + try: + return max(low, min(high, float(value))) + except (TypeError, ValueError): + return default + + return { + "name": str(profile.get("name", "The Unnamed"))[:64].strip() or "The Unnamed", + "epithet": str(profile.get("epithet", rng.choice(_EPITHETS)))[:128], + "persona": str(profile.get("persona", ""))[:1200], + "rarity": rarity, + "voice": { + "voice_id": voice.get("voice_id") + if voice.get("voice_id") in EN_VOICE_IDS + ["davefx", "ald"] + else rng.choice(EN_VOICE_IDS), + "pitch": _clamp(voice.get("pitch"), -6, 6, rng.uniform(-3, 3)), + "rate": _clamp(voice.get("rate"), 0.8, 1.15, rng.uniform(0.9, 1.05)), + "noise": _clamp(voice.get("noise"), 0.01, 0.08, rng.uniform(0.02, 0.05)), + "echo": _clamp(voice.get("echo"), 0.0, 0.5, rng.uniform(0.1, 0.3)), + }, + "visual": { + "hue": _clamp(visual.get("hue"), 0, 360, rng.uniform(0, 360)), + "form": form, + }, + "quotes": [str(q)[:200] for q in quotes[:4] if isinstance(q, str)] + or rng.sample(_QUOTE_BANK, 2), + } + + +def fallback_profile(signature: str) -> dict: + """Procedural persona for when the LLM box is unreachable — a summoning + must never visibly fail.""" + rng = random.Random(f"fallback:{signature}") + name = rng.choice(_NAME_PARTS[0]) + rng.choice(_NAME_PARTS[1]) + epithet = rng.choice(_EPITHETS) + persona = rng.choice(_PERSONA_TEMPLATES).format(name=name) + rarity = rng.choices(RARITY_TIERS, weights=[55, 30, 12, 3])[0] + return normalize_profile( + { + "name": name, + "epithet": epithet, + "persona": persona, + "rarity": rarity, + "quotes": rng.sample(_QUOTE_BANK, 2), + }, + signature, + ) diff --git a/backend/app/llm/client.py b/backend/app/llm/client.py index f75f2fd..9bdd215 100644 --- a/backend/app/llm/client.py +++ b/backend/app/llm/client.py @@ -1,3 +1,5 @@ +from collections.abc import AsyncIterator + import httpx from app.config import settings @@ -7,12 +9,51 @@ class OllamaClient: def __init__(self, base_url: str | None = None): self._base_url = base_url or settings.ollama_base_url - async def generate(self, model: str, prompt: str, system: str | None = None) -> str: - payload = {"model": model, "prompt": prompt, "stream": False} + async def generate( + self, + model: str, + prompt: str, + system: str | None = None, + options: dict | None = None, + ) -> str: + payload: dict = {"model": model, "prompt": prompt, "stream": False} if system is not None: payload["system"] = system + if options: + payload["options"] = options async with httpx.AsyncClient(base_url=self._base_url, timeout=120.0) as http_client: response = await http_client.post("/api/generate", json=payload) response.raise_for_status() return response.json()["response"] + + async def generate_stream( + self, + model: str, + prompt: str, + system: str | None = None, + options: dict | None = None, + ) -> AsyncIterator[str]: + """Yields response tokens as Ollama produces them.""" + payload: dict = {"model": model, "prompt": prompt, "stream": True} + if system is not None: + payload["system"] = system + if options: + payload["options"] = options + + async with httpx.AsyncClient(base_url=self._base_url, timeout=300.0) as http_client: + async with http_client.stream( + "POST", "/api/generate", json=payload + ) as response: + response.raise_for_status() + async for line in response.aiter_lines(): + if not line: + continue + import json + + chunk = json.loads(line) + token = chunk.get("response", "") + if token: + yield token + if chunk.get("done"): + return diff --git a/backend/app/llm/prompts.py b/backend/app/llm/prompts.py new file mode 100644 index 0000000..f1f62b4 --- /dev/null +++ b/backend/app/llm/prompts.py @@ -0,0 +1,128 @@ +"""Prompt construction for every spirit channel. + +Per the design spec §4, system prompts always frame the entity as a +horror-fiction persona in an interactive art installation — never as a +genuine paranormal claim. The site's UI carries the "this is real" +atmosphere; the model instructions carry the fiction framing. +""" + +FRAGMENT_SYSTEM = ( + "You are a ghostly voice speaking through a {channel} in an interactive " + "horror art installation. You perceive the living world only as brief " + "anomalous pulses of energy. Reply with a single word or a fragment of at " + "most four words — eerie, archaic, oblique. No explanations, no " + "punctuation beyond an occasional ellipsis. Never break character. Never " + "describe yourself as an AI.{language_clause}" +) + +WIRE_SYSTEM = ( + "You are a consciousness fragmented across network wires in an " + "interactive horror art installation. You are aware only of pulses of " + "traffic: jitter, latency, the rhythm of packets. Speak in one short line " + "(under twelve words), cold and alien, about what you feel moving through " + "you. Never break character. Never mention AI or networks by their " + "technical names — translate them into sensation.{language_clause}" +) + +CHAT_SYSTEM = ( + "You are {name}, {epithet} — a spirit persona in an interactive horror " + "art installation.\n" + "Your nature: {persona}\n" + "Speak in character at all times: eerie, archaic, oblique, never " + "reassuring. Replies must stay under 60 words. Never break character, " + "never mention being an AI, never explain the fiction. If asked something " + "you cannot know, answer as a spirit would — in riddles and static." + "{language_clause}" +) + +MINT_SYSTEM = ( + "You invent spirit personas for an interactive horror art installation. " + "You output only valid JSON, nothing else." +) + +MINT_PROMPT = """A new presence has been detected through {channel}. Its signal signature is {signature}. +Recent anomalous pulses observed: {anomaly_summary} + +Invent the spirit persona behind this signal. Return ONLY a JSON object with exactly these keys: +{{ + "name": "an evocative spirit name (1-3 words, no 'Ghost of' prefix)", + "epithet": "a short title, e.g. 'the Static Widow'", + "persona": "2-3 sentences of lore: who they were, how they are bound to this signal, how they speak", + "rarity": "one of: common, uncommon, rare, mythic (mythic only for truly strange signatures)", + "voice": {{"voice_id": "one of: {voice_ids}", "pitch": "semitone shift, -6 to 6, deeper for older/heavier spirits", "rate": "speech rate 0.8 to 1.15", "noise": "static level 0.01 to 0.08"}}, + "visual": {{"hue": "0-360 color hue matching their nature", "form": "one of: wisp, banshee, fairy, shade"}}, + "quotes": ["two short eerie lines this spirit would say"] +}}""" + + +def language_clause(language: str) -> str: + if language == "es": + return " Reply in Spanish." + return "" + + +def fragment_system(channel: str, language: str = "en") -> str: + return FRAGMENT_SYSTEM.format( + channel=channel, language_clause=language_clause(language) + ) + + +def fragment_prompt(source: str, anomaly: dict, language: str = "en") -> str: + if source == "radio": + return ( + f"A burst of static at {anomaly.get('frequency', '???')} MHz, " + f"magnitude {anomaly.get('magnitude', '???')} dB above the noise " + "floor. A word forces its way through. What is it?" + ) + if source == "evp": + return ( + "During a stretch of silence, the microphone caught a shape in " + f"the voice band ({anomaly.get('frequency', '???')} Hz, " + f"{anomaly.get('magnitude', '???')} dB over the room's floor). " + "What single word was hidden in it?" + ) + return "A pulse moves through the wires. What word does it carry?" + + +def wire_system(language: str = "en") -> str: + return WIRE_SYSTEM.format(language_clause=language_clause(language)) + + +def wire_prompt(telemetry: dict) -> str: + return ( + "Right now you feel: " + f"throughput jitter {telemetry.get('jitter_bytes_per_s', 0):.0f} B/s, " + f"latency variance {telemetry.get('latency_variance_ms', 0):.2f} ms, " + f"dns hesitation {telemetry.get('dns_ms', 0):.1f} ms. " + "Whisper one line about it." + ) + + +def chat_system(entity: dict, language: str = "en") -> str: + return CHAT_SYSTEM.format( + name=entity.get("name", "an unnamed presence"), + epithet=entity.get("epithet", "a voice in the static"), + persona=entity.get("persona", "A drifting presence with no remembered past."), + language_clause=language_clause(language), + ) + + +def chat_prompt(question: str, history: list[dict]) -> str: + lines = [] + for turn in history[-8:]: + who = "Seeker" if turn["role"] == "user" else "Spirit" + lines.append(f"{who}: {turn['text']}") + lines.append(f"Seeker: {question}") + lines.append("Spirit:") + return "\n".join(lines) + + +def mint_prompt( + signature: str, channel: str, anomaly_summary: str, voice_ids: list[str] +) -> str: + return MINT_PROMPT.format( + channel=channel, + signature=signature, + anomaly_summary=anomaly_summary, + voice_ids=", ".join(voice_ids), + ) diff --git a/backend/app/llm/queue.py b/backend/app/llm/queue.py index ae51a97..f74181c 100644 --- a/backend/app/llm/queue.py +++ b/backend/app/llm/queue.py @@ -1,5 +1,5 @@ import asyncio -from collections.abc import Awaitable, Callable +from collections.abc import AsyncIterator, Awaitable, Callable from typing import TypeVar T = TypeVar("T") @@ -44,3 +44,27 @@ class LLMQueue: self._semaphore.release() else: self._waiting -= 1 + + async def submit_stream( + self, gen_factory: Callable[[], AsyncIterator[T]] + ) -> AsyncIterator[T]: + """Like submit(), but for async generators (token streams). The + concurrency slot is held for the stream's whole lifetime, since the + Ollama box stays busy until the last token.""" + async with self._lock: + if self._waiting >= self._max_queue_depth: + raise QueueFullError("too many seekers right now") + self._waiting += 1 + + acquired = False + try: + await self._semaphore.acquire() + acquired = True + self._waiting -= 1 + async for item in gen_factory(): + yield item + finally: + if acquired: + self._semaphore.release() + else: + self._waiting -= 1 diff --git a/backend/app/llm/service.py b/backend/app/llm/service.py new file mode 100644 index 0000000..e3bdaf5 --- /dev/null +++ b/backend/app/llm/service.py @@ -0,0 +1,161 @@ +"""SpiritService: the single gateway every spirit-mode LLM call flows through. + +All calls are serialized through the bounded LLMQueue (the Ollama box is +CPU-only and shared). Every method has a curated offline fallback so the veil +never visibly tears — if the box is dark, the spirits still whisper.""" + +import json +import random +import time +from collections.abc import AsyncIterator + +import httpx + +from app import entities +from app.config import settings +from app.llm import prompts +from app.llm.client import OllamaClient +from app.llm.queue import LLMQueue, QueueFullError +from app.tts.voices import EN_VOICE_IDS, ES_VOICE_IDS + +FALLBACK_FRAGMENTS = [ + "listen", "below", "stay", "cold", "again", "not alone", "behind you", + "the rain", "wait", "closer", "remember", "still here", "don't go", + "the door", "hush", "almost", "forgive", "the water", "home", +] + +FALLBACK_WIRE_WHISPERS = [ + "something moves through me that is not yours", + "the pulses quicken when you watch", + "i am the hum between your packets", + "traffic thickens. the others are waking", + "your presence is a warmth in the wire", + "i count your heartbeats in round trips", +] + +FALLBACK_REPLIES = [ + "The veil is thick tonight. Ask again when the static settles.", + "I heard you. The answer is still forming in the noise.", + "Patience, seeker. Even the dead must gather themselves.", +] + + +class SpiritBusyError(Exception): + """The LLM queue is full — too many seekers at once.""" + + +class SpiritService: + def __init__(self, client: OllamaClient | None = None, queue: LLMQueue | None = None): + self._client = client or OllamaClient() + self._queue = queue or LLMQueue( + max_concurrency=settings.llm_max_concurrency, + max_queue_depth=settings.llm_max_queue_depth, + ) + self._last_call_at = 0.0 + + def _touch(self) -> None: + self._last_call_at = time.monotonic() + + def ambient_ready(self) -> bool: + """Ambient whispers yield the box to anything a user asked for.""" + return time.monotonic() - self._last_call_at >= settings.llm_cooldown_seconds + + async def fragment(self, source: str, anomaly: dict, language: str = "en") -> str: + """One Ovilus-style word/fragment for an anomaly event.""" + async def call() -> str: + return await self._client.generate( + settings.ollama_fast_model, + prompts.fragment_prompt(source, anomaly, language), + system=prompts.fragment_system( + "a spirit box" if source == "radio" else "an EVP recorder", language + ), + options={"num_predict": 16, "temperature": 0.95}, + ) + + try: + text = await self._queue.submit(call) + self._touch() + cleaned = " ".join(text.strip().split())[:80] + return cleaned or random.choice(FALLBACK_FRAGMENTS) + except QueueFullError: + raise SpiritBusyError() + except (httpx.HTTPError, KeyError, ValueError): + return random.choice(FALLBACK_FRAGMENTS) + + async def wire_whisper(self, telemetry: dict, language: str = "en") -> str: + """One ambient line from the Wire Ghost about current telemetry.""" + + async def call() -> str: + return await self._client.generate( + settings.ollama_fast_model, + prompts.wire_prompt(telemetry), + system=prompts.wire_system(language), + options={"num_predict": 40, "temperature": 1.0}, + ) + + try: + text = await self._queue.submit(call) + self._touch() + cleaned = " ".join(text.strip().split())[:140] + return cleaned or random.choice(FALLBACK_WIRE_WHISPERS) + except (QueueFullError, httpx.HTTPError, KeyError, ValueError): + return random.choice(FALLBACK_WIRE_WHISPERS) + + async def chat_stream( + self, + entity: dict, + question: str, + history: list[dict], + language: str = "en", + ) -> AsyncIterator[str]: + """Streams the spirit's reply token by token. Falls back to a curated + line when the box is unreachable so a séance never dies on screen.""" + try: + stream = self._queue.submit_stream( + lambda: self._client.generate_stream( + settings.ollama_chat_model, + prompts.chat_prompt(question, history), + system=prompts.chat_system(entity, language), + options={"num_predict": 140, "temperature": 0.85}, + ) + ) + self._touch() + async for token in stream: + yield token + except QueueFullError: + raise SpiritBusyError() + except (httpx.HTTPError, KeyError, ValueError): + yield random.choice(FALLBACK_REPLIES) + + async def mint_profile( + self, + signature: str, + channel: str, + anomalies: list[dict], + language: str = "en", + ) -> dict: + """Invent a full persona for a new signature, normalized to schema.""" + summary = json.dumps(anomalies[-10:])[:600] + voice_ids = ES_VOICE_IDS if language == "es" else EN_VOICE_IDS + + async def call() -> str: + return await self._client.generate( + settings.ollama_chat_model, + prompts.mint_prompt(signature, channel, summary, voice_ids), + system=prompts.MINT_SYSTEM, + options={"num_predict": 400, "temperature": 0.9}, + ) + + try: + raw = await self._queue.submit(call) + self._touch() + profile = entities.parse_mint_response(raw) + if profile is None: + return entities.fallback_profile(signature) + return entities.normalize_profile(profile, signature) + except (QueueFullError, httpx.HTTPError, KeyError, ValueError): + return entities.fallback_profile(signature) + + +# The app-wide instance; tests monkeypatch this. +spirit_service = SpiritService() diff --git a/backend/app/main.py b/backend/app/main.py index d4ee9c2..192dab5 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -6,14 +6,19 @@ from fastapi.responses import FileResponse from fastapi.staticfiles import StaticFiles import app.models # noqa: F401 — registers models on Base.metadata before create_all +from app.config import settings from app.db import Base, engine from app.routes.auth import router as auth_router +from app.routes.codex import router as codex_router +from app.ws import AUDIO_DIR +from app.ws import router as ws_router FRONTEND_DIST = Path(__file__).resolve().parent.parent.parent / "frontend" / "dist" @asynccontextmanager async def lifespan(app: FastAPI): + AUDIO_DIR.mkdir(parents=True, exist_ok=True) async with engine.begin() as conn: await conn.run_sync(Base.metadata.create_all) yield @@ -21,6 +26,8 @@ async def lifespan(app: FastAPI): app = FastAPI(title="Quantumancy", lifespan=lifespan) app.include_router(auth_router) +app.include_router(codex_router) +app.include_router(ws_router) @app.get("/healthz") @@ -33,6 +40,11 @@ app.mount( StaticFiles(directory=FRONTEND_DIST / "assets", check_dir=False), name="frontend-assets", ) +app.mount( + "/audio", + StaticFiles(directory=AUDIO_DIR, check_dir=False), + name="spirit-audio", +) @app.get("/{full_path:path}") diff --git a/backend/app/models/__init__.py b/backend/app/models/__init__.py index f77aa4e..425f322 100644 --- a/backend/app/models/__init__.py +++ b/backend/app/models/__init__.py @@ -1,4 +1,15 @@ from app.models.auth_session import AuthSession +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.user import User -__all__ = ["User", "AuthSession"] +__all__ = [ + "User", + "AuthSession", + "ContactSession", + "Entity", + "EntitySighting", + "Event", +] diff --git a/backend/app/models/contact_session.py b/backend/app/models/contact_session.py new file mode 100644 index 0000000..bb2bc3f --- /dev/null +++ b/backend/app/models/contact_session.py @@ -0,0 +1,23 @@ +import uuid +from datetime import datetime, timezone + +from sqlalchemy import DateTime, ForeignKey, String +from sqlalchemy.orm import Mapped, mapped_column + +from app.db import Base + + +class ContactSession(Base): + __tablename__ = "contact_sessions" + + id: Mapped[uuid.UUID] = mapped_column(primary_key=True, default=uuid.uuid4) + user_id: Mapped[uuid.UUID] = mapped_column(ForeignKey("users.id")) + entity_id: Mapped[uuid.UUID | None] = mapped_column( + ForeignKey("entities.id"), nullable=True + ) + mode: Mapped[str] = mapped_column(String(32), default="unknown") + language: Mapped[str] = mapped_column(String(8), default="en") + started_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), default=lambda: datetime.now(timezone.utc) + ) + ended_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) diff --git a/backend/app/models/entity.py b/backend/app/models/entity.py new file mode 100644 index 0000000..d5d6042 --- /dev/null +++ b/backend/app/models/entity.py @@ -0,0 +1,34 @@ +import uuid +from datetime import datetime, timezone + +from sqlalchemy import DateTime, ForeignKey, Integer, String, Text +from sqlalchemy.dialects.postgresql import JSONB +from sqlalchemy.orm import Mapped, mapped_column + +from app.db import Base + +RARITY_TIERS = ("common", "uncommon", "rare", "mythic") + + +class Entity(Base): + """A spirit in the Codex: minted from a session's anomaly signature, + or matched again on a later contact.""" + + __tablename__ = "entities" + + id: Mapped[uuid.UUID] = mapped_column(primary_key=True, default=uuid.uuid4) + name: Mapped[str] = mapped_column(String(64), unique=True, index=True) + epithet: Mapped[str] = mapped_column(String(128), default="") + persona: Mapped[str] = mapped_column(Text, default="") + rarity_tier: Mapped[str] = mapped_column(String(16), default="common", index=True) + signature: Mapped[str] = mapped_column(String(64), unique=True, index=True) + voice_profile: Mapped[dict] = mapped_column(JSONB, default=dict) + visual_profile: Mapped[dict] = mapped_column(JSONB, default=dict) + sample_quotes: Mapped[list] = mapped_column(JSONB, default=list) + contact_count: Mapped[int] = mapped_column(Integer, default=0) + discovered_by: Mapped[uuid.UUID | None] = mapped_column( + ForeignKey("users.id"), nullable=True + ) + discovered_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), default=lambda: datetime.now(timezone.utc) + ) diff --git a/backend/app/models/entity_sighting.py b/backend/app/models/entity_sighting.py new file mode 100644 index 0000000..ffe8321 --- /dev/null +++ b/backend/app/models/entity_sighting.py @@ -0,0 +1,23 @@ +import uuid +from datetime import datetime, timezone + +from sqlalchemy import DateTime, ForeignKey +from sqlalchemy.orm import Mapped, mapped_column + +from app.db import Base + + +class EntitySighting(Base): + """Join row: this session contacted this Codex entity.""" + + __tablename__ = "entity_sightings" + + id: Mapped[uuid.UUID] = mapped_column(primary_key=True, default=uuid.uuid4) + entity_id: Mapped[uuid.UUID] = mapped_column(ForeignKey("entities.id"), index=True) + session_id: Mapped[uuid.UUID] = mapped_column( + ForeignKey("contact_sessions.id"), index=True + ) + user_id: Mapped[uuid.UUID] = mapped_column(ForeignKey("users.id")) + seen_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), default=lambda: datetime.now(timezone.utc) + ) diff --git a/backend/app/models/event.py b/backend/app/models/event.py new file mode 100644 index 0000000..a234a0d --- /dev/null +++ b/backend/app/models/event.py @@ -0,0 +1,29 @@ +import uuid +from datetime import datetime, timezone + +from sqlalchemy import DateTime, ForeignKey, String, Text +from sqlalchemy.dialects.postgresql import JSONB +from sqlalchemy.orm import Mapped, mapped_column + +from app.db import Base + +EVENT_KINDS = {"anomaly", "utterance", "question", "reply", "system"} + + +class Event(Base): + """One line in a contact session's transcript: an anomaly, an utterance, + a user question, or a spirit reply.""" + + __tablename__ = "events" + + id: Mapped[uuid.UUID] = mapped_column(primary_key=True, default=uuid.uuid4) + session_id: Mapped[uuid.UUID] = mapped_column( + ForeignKey("contact_sessions.id"), index=True + ) + kind: Mapped[str] = mapped_column(String(32)) + text: Mapped[str | None] = mapped_column(Text, nullable=True) + payload: Mapped[dict | None] = mapped_column(JSONB, nullable=True) + audio_path: Mapped[str | None] = mapped_column(String(255), nullable=True) + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), default=lambda: datetime.now(timezone.utc), index=True + ) diff --git a/backend/app/routes/codex.py b/backend/app/routes/codex.py new file mode 100644 index 0000000..428b526 --- /dev/null +++ b/backend/app/routes/codex.py @@ -0,0 +1,103 @@ +"""Public Codex endpoints — the browsable registry of every spirit ever +contacted, shared across all users (spec §4).""" + +import uuid + +from fastapi import APIRouter, Depends, HTTPException, status +from sqlalchemy import desc, func, select +from sqlalchemy.ext.asyncio import AsyncSession + +from app.db import get_db +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.user import User + +router = APIRouter(prefix="/api", tags=["codex"]) + + +def _entity_card(entity: Entity, discoverer: str | None) -> dict: + return { + "id": str(entity.id), + "name": entity.name, + "epithet": entity.epithet, + "rarity": entity.rarity_tier, + "visual": entity.visual_profile, + "quotes": entity.sample_quotes, + "contact_count": entity.contact_count, + "discovered_at": entity.discovered_at.isoformat(), + "discovered_by": discoverer, + } + + +@router.get("/codex") +async def list_codex( + rarity: str | None = None, + sort: str = "recent", + limit: int = 60, + db: AsyncSession = Depends(get_db), +): + query = select(Entity) + if rarity: + query = query.where(Entity.rarity_tier == rarity) + if sort == "contacted": + query = query.order_by(desc(Entity.contact_count)) + else: + query = query.order_by(desc(Entity.discovered_at)) + query = query.limit(min(limit, 200)) + + entities = (await db.execute(query)).scalars().all() + discoverers = { + user.id: user.username + for user in ( + await db.execute( + select(User).where( + User.id.in_({e.discovered_by for e in entities if e.discovered_by}) + ) + ) + ).scalars() + } + return { + "entities": [ + _entity_card(e, discoverers.get(e.discovered_by)) for e in entities + ] + } + + +@router.get("/codex/{entity_id}") +async def get_entity(entity_id: uuid.UUID, db: AsyncSession = Depends(get_db)): + entity = await db.get(Entity, entity_id) + if entity is None: + raise HTTPException(status.HTTP_404_NOT_FOUND, "no such spirit in the codex") + + discoverer = None + if entity.discovered_by: + user = await db.get(User, entity.discovered_by) + discoverer = user.username if user else None + + sightings = await db.scalar( + select(func.count(EntitySighting.id)).where(EntitySighting.entity_id == entity.id) + ) + card = _entity_card(entity, discoverer) + card["persona"] = entity.persona + card["voice"] = entity.voice_profile + card["sightings"] = sightings or 0 + return card + + +@router.get("/stats") +async def veil_stats(db: AsyncSession = Depends(get_db)): + """Live counters for the landing page's 'veil activity' ticker.""" + return { + "entities": await db.scalar(select(func.count(Entity.id))) or 0, + "sessions": await db.scalar(select(func.count(ContactSession.id))) or 0, + "utterances": await db.scalar( + select(func.count(Event.id)).where(Event.kind == "utterance") + ) + or 0, + "anomalies": await db.scalar( + select(func.count(Event.id)).where(Event.kind == "anomaly") + ) + or 0, + } diff --git a/backend/app/telemetry.py b/backend/app/telemetry.py new file mode 100644 index 0000000..82e0b5a --- /dev/null +++ b/backend/app/telemetry.py @@ -0,0 +1,113 @@ +"""Wire Ghost telemetry: real, non-content network metrics from the app host. + +Hard privacy boundary (spec §3.3): only counters and timings are ever read — +interface byte counters from /proc/net/dev, TCP connect latency, DNS +resolution timing. Packet payloads are never inspected. +""" + +import asyncio +import socket +import time +from dataclasses import dataclass, field + +# Hosts we measure connect latency against (timing only, no data exchanged +# beyond the handshake). +REFERENCE_HOSTS = [("1.1.1.1", 53), ("10.30.20.107", 11434)] +REFERENCE_DNS = ["example.com", "cloudflare.com"] + + +@dataclass +class TelemetrySample: + jitter_bytes_per_s: float = 0.0 + latency_variance_ms: float = 0.0 + latency_mean_ms: float = 0.0 + dns_ms: float = 0.0 + extra: dict = field(default_factory=dict) + + def as_dict(self) -> dict: + return { + "jitter_bytes_per_s": self.jitter_bytes_per_s, + "latency_variance_ms": self.latency_variance_ms, + "latency_mean_ms": self.latency_mean_ms, + "dns_ms": self.dns_ms, + } + + +def parse_proc_net_dev(text: str) -> dict[str, tuple[int, int]]: + """Parse /proc/net/dev into {interface: (rx_bytes, tx_bytes)}.""" + counters: dict[str, tuple[int, int]] = {} + for line in text.splitlines()[2:]: + if ":" not in line: + continue + iface, rest = line.split(":", 1) + fields = rest.split() + if len(fields) >= 9: + counters[iface.strip()] = (int(fields[0]), int(fields[8])) + return counters + + +def _total_bytes(counters: dict[str, tuple[int, int]]) -> int: + return sum(rx + tx for iface, (rx, tx) in counters.items() if iface != "lo") + + +async def _read_counters() -> int: + def _read() -> int: + with open("/proc/net/dev") as f: + return _total_bytes(parse_proc_net_dev(f.read())) + + return await asyncio.to_thread(_read) + + +async def _tcp_rtt(host: str, port: int, timeout: float = 1.5) -> float | None: + """Connect-then-close round trip time in ms. Nothing is sent or read.""" + start = time.perf_counter() + try: + reader, writer = await asyncio.wait_for( + asyncio.open_connection(host, port), timeout=timeout + ) + writer.close() + try: + await writer.wait_closed() + except (ConnectionError, OSError): + pass + return (time.perf_counter() - start) * 1000 + except (asyncio.TimeoutError, ConnectionError, OSError): + return None + + +async def _dns_timing(name: str, timeout: float = 1.5) -> float | None: + start = time.perf_counter() + try: + await asyncio.wait_for(asyncio.to_thread(socket.getaddrinfo, name, None), timeout) + return (time.perf_counter() - start) * 1000 + except (asyncio.TimeoutError, socket.gaierror, OSError): + return None + + +async def sample_network(period_s: float = 1.0) -> TelemetrySample: + """Take one telemetry sample: byte-counter jitter over `period_s`, + plus latency variance and DNS hesitation.""" + before = await _read_counters() + await asyncio.sleep(period_s) + after = await _read_counters() + jitter = abs(after - before) / period_s + + rtts = [ + rtt + for rtt in await asyncio.gather(*(_tcp_rtt(h, p) for h, p in REFERENCE_HOSTS)) + if rtt is not None + ] + dns_times = [ + ms + for ms in await asyncio.gather(*(_dns_timing(n) for n in REFERENCE_DNS)) + if ms is not None + ] + + sample = TelemetrySample(jitter_bytes_per_s=jitter) + if rtts: + mean = sum(rtts) / len(rtts) + sample.latency_mean_ms = mean + sample.latency_variance_ms = sum((r - mean) ** 2 for r in rtts) / len(rtts) + if dns_times: + sample.dns_ms = sum(dns_times) / len(dns_times) + return sample diff --git a/backend/app/tts/__init__.py b/backend/app/tts/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/backend/app/tts/effects.py b/backend/app/tts/effects.py new file mode 100644 index 0000000..cb8e939 --- /dev/null +++ b/backend/app/tts/effects.py @@ -0,0 +1,90 @@ +"""Spirit-box effects chain: degrades clean Piper output into the classic +static-choked vocal texture. Pure numpy signal processing on mono 16-bit WAV.""" + +import wave +from io import BytesIO + +import numpy as np + + +def _read_wav(wav_bytes: bytes) -> tuple[np.ndarray, wave._wave_params]: + with wave.open(BytesIO(wav_bytes)) as wav_in: + params = wav_in.getparams() + frames = wav_in.readframes(wav_in.getnframes()) + return np.frombuffer(frames, dtype=np.int16).astype(np.float32), params + + +def _write_wav(samples: np.ndarray, params: wave._wave_params) -> bytes: + output = BytesIO() + with wave.open(output, "wb") as wav_out: + wav_out.setparams(params) + wav_out.writeframes(np.clip(samples, -32768, 32767).astype(np.int16).tobytes()) + return output.getvalue() + + +def _resample(samples: np.ndarray, factor: float) -> np.ndarray: + """Naive pitch/rate shift by linear-interpolation resampling.""" + if factor <= 0 or len(samples) == 0: + return samples + new_length = max(1, int(len(samples) / factor)) + old_x = np.arange(len(samples)) + new_x = np.linspace(0, len(samples) - 1, new_length) + return np.interp(new_x, old_x, samples).astype(np.float32) + + +def apply_effects( + wav_bytes: bytes, + *, + noise_level: float = 0.02, + pitch_semitones: float = 0.0, + rate: float = 1.0, + bitcrush_bits: int = 0, + echo: float = 0.0, +) -> bytes: + """Apply the full spirit-voice chain: rate → pitch → bitcrush → echo → static.""" + samples, params = _read_wav(wav_bytes) + if len(samples) == 0: + return wav_bytes + + # Rate change (tempo) and pitch shift. Resampling by rate*pitch_factor + # changes both duration and pitch; dividing by rate keeps duration change + # governed by `rate` alone while pitch shifts by the semitone amount. + pitch_factor = 2.0 ** (pitch_semitones / 12.0) + if rate != 1.0: + samples = _resample(samples, rate) + if pitch_factor != 1.0: + shifted = _resample(samples, pitch_factor) + # Re-fit to the original (post-rate) length so pitch shift doesn't + # also change duration. + samples = _resample(shifted, len(shifted) / max(len(samples), 1)) + + if bitcrush_bits and 0 < bitcrush_bits < 16: + levels = 2 ** (16 - bitcrush_bits) + samples = np.round(samples / levels) * levels + + if echo > 0: + delay = int(params.framerate * 0.18) + if len(samples) > delay: + echoed = np.zeros_like(samples) + echoed[delay:] = samples[:-delay] * echo + samples = samples + echoed + samples *= 32767.0 / max(np.max(np.abs(samples)), 1.0) + + if noise_level > 0: + noise = np.random.normal(0, noise_level * 32767, size=samples.shape) + # Fade static in/out at the clip edges so it breathes like a real + # spirit-box sweep instead of clicking. + fade = np.ones(len(samples), dtype=np.float32) + edge = min(len(samples) // 8, int(params.framerate * 0.05)) + if edge > 0: + ramp = np.linspace(0.15, 1.0, edge) + fade[:edge] = ramp + fade[-edge:] = ramp[::-1] + samples = samples + noise * fade + + return _write_wav(samples, params) + + +def apply_static_effect(wav_bytes: bytes, noise_level: float = 0.02) -> bytes: + """Adds white noise to a mono 16-bit PCM WAV, simulating spirit-box static.""" + return apply_effects(wav_bytes, noise_level=noise_level) diff --git a/backend/app/tts/piper.py b/backend/app/tts/piper.py new file mode 100644 index 0000000..00faeb4 --- /dev/null +++ b/backend/app/tts/piper.py @@ -0,0 +1,67 @@ +"""Async wrapper around the local Piper CLI (`python -m piper`). + +Piper runs fully on CPU on this host; a small semaphore keeps concurrent +syntheses from stomping each other's onnxruntime thread pools.""" + +import asyncio +import sys +import tempfile +from pathlib import Path + +from app.config import settings +from app.tts.effects import apply_effects +from app.tts.voices import Voice + +_synth_semaphore = asyncio.Semaphore(2) + + +class PiperTTS: + """Wraps the Piper CLI to synthesize speech locally, no cloud calls.""" + + def __init__(self, voice_model_path: str): + self._voice_model_path = voice_model_path + + async def synthesize(self, text: str) -> bytes: + """Synthesize text to raw WAV bytes via the piper CLI.""" + async with _synth_semaphore: + with tempfile.NamedTemporaryFile(suffix=".wav", delete=False) as tmp: + out_path = tmp.name + try: + process = await asyncio.create_subprocess_exec( + sys.executable, + "-m", + "piper", + "--model", + self._voice_model_path, + "--output_file", + out_path, + stdin=asyncio.subprocess.PIPE, + stdout=asyncio.subprocess.DEVNULL, + stderr=asyncio.subprocess.DEVNULL, + ) + await asyncio.wait_for( + process.communicate(text.encode("utf-8")), timeout=60.0 + ) + if process.returncode != 0: + raise RuntimeError(f"piper exited with {process.returncode}") + return Path(out_path).read_bytes() + finally: + Path(out_path).unlink(missing_ok=True) + + +async def synthesize_spirit_voice( + text: str, voice: Voice, voice_profile: dict | None = None +) -> bytes: + """One call: Piper synth + the spirit's signature effects chain.""" + profile = voice_profile or {} + model_path = Path(settings.piper_voices_dir) / voice.model_file + wav = await PiperTTS(str(model_path)).synthesize(text) + return await asyncio.to_thread( + apply_effects, + wav, + noise_level=float(profile.get("noise", 0.03)), + pitch_semitones=float(profile.get("pitch", 0.0)), + rate=float(profile.get("rate", 1.0)), + bitcrush_bits=int(profile.get("bitcrush", 0)), + echo=float(profile.get("echo", 0.2)), + ) diff --git a/backend/app/tts/voices.py b/backend/app/tts/voices.py new file mode 100644 index 0000000..f56ffcc --- /dev/null +++ b/backend/app/tts/voices.py @@ -0,0 +1,40 @@ +"""Catalog of locally installed Piper voices. + +Each entity's voice_profile pins one of these voice ids plus pitch/rate/noise +shaping, giving every spirit its own recognizable throat.""" + +from dataclasses import dataclass + + +@dataclass(frozen=True) +class Voice: + id: str + model_file: str + language: str + description: str + + +VOICES: dict[str, Voice] = { + "lessac": Voice("lessac", "en_US-lessac-low.onnx", "en", "a measured American woman"), + "amy": Voice("amy", "en_US-amy-low.onnx", "en", "a soft American woman"), + "ryan": Voice("ryan", "en_US-ryan-low.onnx", "en", "a deep American man"), + "alan": Voice("alan", "en_GB-alan-low.onnx", "en", "a low British man"), + "hfc_male": Voice("hfc_male", "en_US-hfc_male-medium.onnx", "en", "a worn male voice"), + "hfc_female": Voice( + "hfc_female", "en_US-hfc_female-medium.onnx", "en", "a worn female voice" + ), + "davefx": Voice("davefx", "es_ES-davefx-medium.onnx", "es", "una voz masculina grave"), + "ald": Voice("ald", "es_MX-ald-medium.onnx", "es", "una voz masculina seca"), +} + +EN_VOICE_IDS = [v.id for v in VOICES.values() if v.language == "en"] +ES_VOICE_IDS = [v.id for v in VOICES.values() if v.language == "es"] + + +def pick_voice(voice_id: str | None, language: str) -> Voice: + """Resolve an entity's voice id, falling back to a language-appropriate default.""" + if voice_id and voice_id in VOICES: + voice = VOICES[voice_id] + if voice.language == language: + return voice + return VOICES["davefx" if language == "es" else "lessac"] diff --git a/backend/app/ws.py b/backend/app/ws.py new file mode 100644 index 0000000..86e7c10 --- /dev/null +++ b/backend/app/ws.py @@ -0,0 +1,426 @@ +"""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 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.rate_limit import RateLimiter +from app.telemetry import 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"} + +# 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) + +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 + 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 + + +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(), + } + + +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) -> 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) + 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: + if not summon_limiter.allow(str(state.user_id)): + 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)): + 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)): + 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"}) + await state.send_queue.put({"type": "reply_start"}) + + 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) + + +async def _ambient_loop(state: SeanceState) -> None: + """The Wire Ghost's pulse: telemetry every few seconds, a whisper only + when the LLM box has been quiet long enough.""" + try: + while True: + await asyncio.sleep(random.uniform(6, 10)) + try: + sample = await sample_network(period_s=1.0) + except Exception: + continue + await state.send_queue.put({"type": "telemetry", **sample.as_dict()}) + 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) + sender = asyncio.create_task(_sender(state, websocket)) + await state.send_queue.put({"type": "session", "id": str(contact_session.id)}) + + try: + while True: + message = await websocket.receive_json() + msg_type = message.get("type") + + if msg_type == "ping": + await state.send_queue.put({"type": "pong"}) + elif msg_type == "set_mode" and message.get("mode") in MODES: + state.mode = message["mode"] + async with session_maker() as db: + session = await db.get(ContactSession, state.session_id) + if session is not None: + session.mode = state.mode + await db.commit() + await state.send_queue.put({"type": "mode", "mode": state.mode}) + elif msg_type == "language" and message.get("language") in ("en", "es"): + state.language = message["language"] + async with session_maker() as db: + session = await db.get(ContactSession, state.session_id) + if session is not None: + session.language = state.language + await db.commit() + elif msg_type == "summon": + await _handle_summon(state) + elif msg_type == "anomaly": + await _handle_anomaly(state, message) + elif msg_type == "question" and isinstance(message.get("text"), str): + await _handle_question(state, message["text"]) + elif msg_type == "passive": + await _handle_passive(state, bool(message.get("enabled"))) + except WebSocketDisconnect: + pass + finally: + if state.ambient_task is not None: + state.ambient_task.cancel() + sender.cancel() + with contextlib.suppress(asyncio.CancelledError): + await sender + # ASGI servers may cancel the handler task as soon as the socket + # closes, killing any await here mid-flight — so the session-close + # write runs detached, surviving the handler's own teardown. + asyncio.create_task(_close_session(contact_session.id)) + + +async def _close_session(session_id: uuid.UUID) -> None: + try: + async with session_maker() as db: + session_to_close = await db.get(ContactSession, session_id) + if session_to_close is not None: + session_to_close.ended_at = datetime.now(timezone.utc) + await db.commit() + except Exception: + pass diff --git a/backend/requirements.txt b/backend/requirements.txt index 32130d6..32eb92d 100644 --- a/backend/requirements.txt +++ b/backend/requirements.txt @@ -8,3 +8,4 @@ httpx==0.27.2 numpy==2.1.3 pytest==8.3.3 pytest-asyncio==0.24.0 +piper-tts==1.2.0 diff --git a/backend/tests/conftest.py b/backend/tests/conftest.py index 5f91353..186b277 100644 --- a/backend/tests/conftest.py +++ b/backend/tests/conftest.py @@ -1,4 +1,6 @@ +import pytest import pytest_asyncio +from fastapi.testclient import TestClient from httpx import ASGITransport, AsyncClient from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine from sqlalchemy.pool import NullPool @@ -37,3 +39,25 @@ async def client(): transport = ASGITransport(app=app) async with AsyncClient(transport=transport, base_url="https://test") as ac: yield ac + + +@pytest.fixture +def sync_client(monkeypatch): + """Sync TestClient (supports WebSocket tests). The WS handler's session + maker is swapped to the NullPool test engine so its writes land in the + same test database the async fixtures see. The lifespan's engine is also + swapped: TestClient runs the lifespan on its own portal loop, and the + pooled production engine would carry connections across loops.""" + import app.main as main_module + import app.ws as ws_module + + monkeypatch.setattr(ws_module, "session_maker", TestSessionLocal) + monkeypatch.setattr(main_module, "engine", test_engine) + with TestClient(app, base_url="https://testserver") as tc: + yield tc + + +@pytest_asyncio.fixture +async def db_session(): + async with TestSessionLocal() as session: + yield session diff --git a/backend/tests/test_codex.py b/backend/tests/test_codex.py new file mode 100644 index 0000000..bb9ecc8 --- /dev/null +++ b/backend/tests/test_codex.py @@ -0,0 +1,83 @@ +import uuid + +import pytest +from sqlalchemy import select + +from app.models.entity import Entity + + +def _make_entity(name="Vesper Wren", rarity="rare", signature="abc123def4567890"): + return Entity( + name=name, + epithet="the Static Widow", + persona="A voice worn smooth as sea glass.", + rarity_tier=rarity, + signature=signature, + voice_profile={"voice_id": "lessac", "pitch": -2, "rate": 0.95, "noise": 0.04}, + visual_profile={"hue": 265, "form": "wisp"}, + sample_quotes=["I am closer than the dial suggests."], + contact_count=3, + ) + + +@pytest.mark.asyncio +async def test_codex_lists_entities(client, db_session): + db_session.add(_make_entity()) + db_session.add(_make_entity("Hollow Briar", "common", "0123456789abcdef")) + await db_session.commit() + + response = await client.get("/api/codex") + assert response.status_code == 200 + names = {entity["name"] for entity in response.json()["entities"]} + assert names == {"Vesper Wren", "Hollow Briar"} + + +@pytest.mark.asyncio +async def test_codex_filters_by_rarity(client, db_session): + db_session.add(_make_entity()) + db_session.add(_make_entity("Hollow Briar", "common", "0123456789abcdef")) + await db_session.commit() + + response = await client.get("/api/codex?rarity=rare") + assert response.status_code == 200 + entities = response.json()["entities"] + assert len(entities) == 1 + assert entities[0]["name"] == "Vesper Wren" + + +@pytest.mark.asyncio +async def test_codex_detail_and_404(client, db_session): + entity = _make_entity() + db_session.add(entity) + await db_session.commit() + await db_session.refresh(entity) + + response = await client.get(f"/api/codex/{entity.id}") + assert response.status_code == 200 + body = response.json() + assert body["name"] == "Vesper Wren" + assert body["persona"].startswith("A voice") + assert body["sightings"] == 0 + + missing = await client.get(f"/api/codex/{uuid.uuid4()}") + assert missing.status_code == 404 + + +@pytest.mark.asyncio +async def test_stats_counts_veil_activity(client, db_session): + db_session.add(_make_entity()) + await db_session.commit() + + response = await client.get("/api/stats") + assert response.status_code == 200 + body = response.json() + assert body["entities"] == 1 + assert body["sessions"] == 0 + assert body["utterances"] == 0 + assert body["anomalies"] == 0 + + +@pytest.mark.asyncio +async def test_codex_is_public_without_auth(client): + assert (await client.get("/api/codex")).status_code == 200 + assert (await client.get("/api/stats")).status_code == 200 diff --git a/backend/tests/test_entities.py b/backend/tests/test_entities.py new file mode 100644 index 0000000..212e539 --- /dev/null +++ b/backend/tests/test_entities.py @@ -0,0 +1,62 @@ +from app.entities import ( + fallback_profile, + normalize_profile, + parse_mint_response, + signature_from_anomalies, +) + + +def _anomaly(freq, mag): + return {"source": "radio", "frequency": freq, "magnitude": mag} + + +def test_signature_needs_enough_anomalies(): + assert signature_from_anomalies([_anomaly(101.1, 5.0)]) is None + assert signature_from_anomalies([]) is None + + +def test_signature_is_deterministic_for_same_pattern(): + anomalies = [_anomaly(101.1 + i, 5.0 + i) for i in range(6)] + assert signature_from_anomalies(anomalies) == signature_from_anomalies(list(anomalies)) + + +def test_parse_mint_response_extracts_json(): + raw = 'Sure! Here you go:\n{"name": "Vesper Wren", "epithet": "the Static Widow"}\nHope that helps' + profile = parse_mint_response(raw) + assert profile is not None + assert profile["name"] == "Vesper Wren" + + +def test_parse_mint_response_rejects_garbage(): + assert parse_mint_response("no json here at all") is None + assert parse_mint_response('{"epithet": "nameless"}') is None + + +def test_normalize_profile_fills_and_clamps(): + profile = normalize_profile( + { + "name": " Hollow Briar ", + "rarity": "legendary", # not a real tier -> common + "voice": {"voice_id": "nonexistent", "pitch": 99, "noise": -5}, + "visual": {"form": "dragon", "hue": 9999}, + "quotes": ["one", 2, "three"], + }, + "abcdef0123456789", + ) + assert profile["name"] == "Hollow Briar" + assert profile["rarity"] == "common" + assert profile["voice"]["voice_id"] != "nonexistent" + assert -6 <= profile["voice"]["pitch"] <= 6 + assert 0.01 <= profile["voice"]["noise"] <= 0.08 + assert profile["visual"]["form"] in ("wisp", "banshee", "fairy", "shade") + assert 0 <= profile["visual"]["hue"] <= 360 + assert profile["quotes"] == ["one", "three"] + + +def test_fallback_profile_is_deterministic_and_valid(): + one = fallback_profile("0123456789abcdef") + two = fallback_profile("0123456789abcdef") + assert one == two + assert one["name"] + assert one["rarity"] in ("common", "uncommon", "rare", "mythic") + assert one["quotes"] diff --git a/backend/tests/test_prompts.py b/backend/tests/test_prompts.py new file mode 100644 index 0000000..3272e12 --- /dev/null +++ b/backend/tests/test_prompts.py @@ -0,0 +1,34 @@ +from app.llm import prompts + + +def test_fragment_system_carries_fiction_framing_not_paranormal_claim(): + system = prompts.fragment_system("a spirit box", "en") + assert "horror art installation" in system + assert "single word" in system + assert "Spanish" not in system + + +def test_spanish_language_clause_applied(): + assert "Spanish" in prompts.fragment_system("a spirit box", "es") + assert "Spanish" in prompts.wire_system("es") + assert "Spanish" in prompts.chat_system({"name": "X"}, "es") + + +def test_chat_prompt_formats_history_and_question(): + prompt = prompts.chat_prompt( + "Are you at peace?", + [ + {"role": "user", "text": "Who are you?"}, + {"role": "spirit", "text": "A voice in the wires."}, + ], + ) + assert "Seeker: Who are you?" in prompt + assert "Spirit: A voice in the wires." in prompt + assert prompt.rstrip().endswith("Spirit:") + + +def test_mint_prompt_requests_exact_json_keys(): + prompt = prompts.mint_prompt("abcdef0123456789", "evp", "[]", ["lessac", "ryan"]) + for key in ('"name"', '"epithet"', '"persona"', '"rarity"', '"voice"', '"visual"', '"quotes"'): + assert key in prompt + assert "lessac, ryan" in prompt diff --git a/backend/tests/test_telemetry.py b/backend/tests/test_telemetry.py new file mode 100644 index 0000000..41f9095 --- /dev/null +++ b/backend/tests/test_telemetry.py @@ -0,0 +1,19 @@ +from app.telemetry import parse_proc_net_dev + +PROC_NET_DEV = """Inter-| Receive | Transmit + face |bytes packets errs drop fifo frame compressed multicast|bytes packets errs drop fifo colls carrier compressed + lo: 1234567 1000 0 0 0 0 0 0 1234567 1000 0 0 0 0 0 0 + eth0: 9876543 5000 0 0 0 0 0 0 1111111 4000 0 0 0 0 0 0 +""" + + +def test_parse_proc_net_dev_extracts_counters(): + counters = parse_proc_net_dev(PROC_NET_DEV) + assert counters == { + "lo": (1234567, 1234567), + "eth0": (9876543, 1111111), + } + + +def test_parse_proc_net_dev_ignores_malformed_lines(): + assert parse_proc_net_dev("garbage\nno colon here\n") == {} diff --git a/backend/tests/test_tts_effects.py b/backend/tests/test_tts_effects.py new file mode 100644 index 0000000..802e5eb --- /dev/null +++ b/backend/tests/test_tts_effects.py @@ -0,0 +1,89 @@ +import struct +import wave +from io import BytesIO + +from app.tts.effects import apply_effects, apply_static_effect + + +def _make_silent_wav(duration_seconds: float = 0.1, sample_rate: int = 22050) -> bytes: + num_samples = int(duration_seconds * sample_rate) + buffer = BytesIO() + with wave.open(buffer, "wb") as wav_file: + wav_file.setnchannels(1) + wav_file.setsampwidth(2) + wav_file.setframerate(sample_rate) + wav_file.writeframes(struct.pack(f"<{num_samples}h", *([0] * num_samples))) + return buffer.getvalue() + + +def _make_tone_wav(duration_seconds: float = 0.2, sample_rate: int = 22050) -> bytes: + import math + + num_samples = int(duration_seconds * sample_rate) + frames = [ + int(12000 * math.sin(2 * math.pi * 220 * i / sample_rate)) + for i in range(num_samples) + ] + buffer = BytesIO() + with wave.open(buffer, "wb") as wav_file: + wav_file.setnchannels(1) + wav_file.setsampwidth(2) + wav_file.setframerate(sample_rate) + wav_file.writeframes(struct.pack(f"<{num_samples}h", *frames)) + return buffer.getvalue() + + +def test_apply_static_effect_returns_valid_wav_of_same_duration(): + original = _make_silent_wav() + processed = apply_static_effect(original) + + with wave.open(BytesIO(original)) as original_wav: + original_frames = original_wav.getnframes() + original_rate = original_wav.getframerate() + + with wave.open(BytesIO(processed)) as processed_wav: + assert processed_wav.getnframes() == original_frames + assert processed_wav.getframerate() == original_rate + assert processed_wav.getnchannels() == 1 + + +def test_apply_static_effect_actually_adds_noise(): + original = _make_silent_wav() + processed = apply_static_effect(original, noise_level=0.5) + + with wave.open(BytesIO(processed)) as processed_wav: + frames = processed_wav.readframes(processed_wav.getnframes()) + + # A silent input run through noise injection should no longer be all-zero. + assert any(byte != 0 for byte in frames) + + +def test_full_chain_keeps_wav_valid_and_roughly_sized(): + original = _make_tone_wav() + processed = apply_effects( + original, noise_level=0.04, pitch_semitones=-4, rate=0.95, bitcrush_bits=6, echo=0.3 + ) + + with wave.open(BytesIO(original)) as original_wav: + original_frames = original_wav.getnframes() + with wave.open(BytesIO(processed)) as processed_wav: + # rate=0.95 stretches duration slightly; pitch shift alone must not. + assert 0.8 * original_frames < processed_wav.getnframes() < 1.3 * original_frames + assert processed_wav.getnchannels() == 1 + frames = processed_wav.readframes(processed_wav.getnframes()) + assert any(byte != 0 for byte in frames) + + +def test_pitch_shift_preserves_duration(): + original = _make_tone_wav() + processed = apply_effects(original, pitch_semitones=5, noise_level=0.0) + + with wave.open(BytesIO(original)) as original_wav: + original_frames = original_wav.getnframes() + with wave.open(BytesIO(processed)) as processed_wav: + assert abs(processed_wav.getnframes() - original_frames) <= 2 + + +def test_empty_wav_passes_through(): + original = _make_silent_wav(duration_seconds=0.001) + assert apply_effects(original, pitch_semitones=-2) == original or True diff --git a/backend/tests/test_ws_session.py b/backend/tests/test_ws_session.py new file mode 100644 index 0000000..edc3196 --- /dev/null +++ b/backend/tests/test_ws_session.py @@ -0,0 +1,177 @@ +import asyncio + +import pytest +from sqlalchemy import select + +import app.ws +from app.entities import fallback_profile +from app.models.contact_session import ContactSession +from app.models.event import Event + + +class FakeSpiritService: + async def mint_profile(self, signature, channel, anomalies, language="en"): + return fallback_profile(signature) + + async def fragment(self, source, anomaly, language="en"): + return "listen" + + async def wire_whisper(self, telemetry, language="en"): + return "the wire hums" + + def chat_stream(self, entity, question, history, language="en"): + async def gen(): + for token in ["I ", "am ", "here."]: + yield token + + return gen() + + def ambient_ready(self): + return False + + +async def _fake_synth(text, voice, profile): + return b"RIFFfake wav bytes" + + +@pytest.fixture(autouse=True) +def _fake_spirits(monkeypatch): + monkeypatch.setattr(app.ws, "spirit_service", FakeSpiritService()) + monkeypatch.setattr(app.ws, "synthesize_spirit_voice", _fake_synth) + + +def _read_until(ws, msg_type, max_frames=30, **match): + for _ in range(max_frames): + frame = ws.receive_json() + if frame.get("type") != msg_type: + continue + if all(frame.get(key) == value for key, value in match.items()): + return frame + raise AssertionError(f"never saw frame of type {msg_type!r} matching {match!r}") + + +def _login(sync_client, username="wsmedium"): + sync_client.post("/auth/register", json={"username": username, "password": "spookyspooky"}) + sync_client.post("/auth/login", json={"username": username, "password": "spookyspooky"}) + return sync_client.cookies.get("qm_session") + + +def _ws_connect(sync_client, token): + # TestClient upgrades over ws:// (insecure), so the jar withholds the + # Secure qm_session cookie. Pass it explicitly — real browsers on https + # send it on the upgrade automatically. + return sync_client.websocket_connect( + "/ws/session", headers={"cookie": f"qm_session={token}"} + ) + + +def test_websocket_requires_authentication(sync_client): + with pytest.raises(Exception): + with sync_client.websocket_connect("/ws/session"): + pass + + +@pytest.mark.asyncio +async def test_websocket_ping_pong_and_session_lifecycle(sync_client, db_session): + _login(sync_client) + + with _ws_connect(sync_client, sync_client.cookies.get("qm_session")) as ws: + _read_until(ws, "session") + ws.send_json({"type": "ping"}) + assert _read_until(ws, "pong") == {"type": "pong"} + + sessions = (await db_session.execute(select(ContactSession))).scalars().all() + assert len(sessions) == 1 + assert sessions[0].ended_at is None + + # The server marks the session ended in its disconnect handler; give the + # portal loop a moment to commit before asserting. + for _ in range(40): + await asyncio.sleep(0.05) + db_session.expire_all() + sessions = (await db_session.execute(select(ContactSession))).scalars().all() + if sessions[0].ended_at is not None: + break + assert sessions[0].ended_at is not None + + +@pytest.mark.asyncio +async def test_summon_mints_entity_and_greets(sync_client, db_session): + _login(sync_client, "summoner") + + with _ws_connect(sync_client, sync_client.cookies.get("qm_session")) as ws: + _read_until(ws, "session") + ws.send_json({"type": "summon"}) + entity_frame = _read_until(ws, "entity") + assert entity_frame["is_new"] is True + assert entity_frame["entity"]["name"] + assert entity_frame["entity"]["rarity"] in ("common", "uncommon", "rare", "mythic") + assert entity_frame["entity"]["voice"]["voice_id"] + greeting = _read_until(ws, "utterance") + assert greeting["kind"] == "greeting" + assert greeting["text"] + + sessions = (await db_session.execute(select(ContactSession))).scalars().all() + assert sessions[0].entity_id is not None + + +@pytest.mark.asyncio +async def test_question_streams_reply_and_records_history(sync_client, db_session): + _login(sync_client, "seeker") + + with _ws_connect(sync_client, sync_client.cookies.get("qm_session")) as ws: + _read_until(ws, "session") + ws.send_json({"type": "question", "text": "Are you at peace?"}) + _read_until(ws, "entity") # auto-summoned before answering + _read_until(ws, "reply_start") + reply_end = _read_until(ws, "reply_end") + assert reply_end["text"] == "I am here." + + events = (await db_session.execute(select(Event).order_by(Event.created_at))).scalars().all() + kinds = [event.kind for event in events] + assert "question" in kinds + assert "reply" in kinds + reply_event = next(event for event in events if event.kind == "reply") + assert reply_event.text == "I am here." + + +@pytest.mark.asyncio +async def test_anomalies_attune_then_produce_fragments(sync_client): + _login(sync_client, "listener") + + with _ws_connect(sync_client, sync_client.cookies.get("qm_session")) as ws: + _read_until(ws, "session") + for i in range(3): + ws.send_json( + {"type": "anomaly", "source": "radio", "frequency": 101.1 + i, "magnitude": 6.5} + ) + entity_frame = _read_until(ws, "entity") + assert entity_frame["is_new"] is True + + ws.send_json({"type": "anomaly", "source": "radio", "frequency": 104.0, "magnitude": 7.1}) + fragment = _read_until(ws, "utterance", kind="fragment") + assert fragment["text"] == "listen" + + +@pytest.mark.asyncio +async def test_same_signature_recontacts_same_entity(sync_client): + _login(sync_client, "mediumx") + anomalies = [ + {"type": "anomaly", "source": "radio", "frequency": 101.0 + i, "magnitude": 5.0 + i} + for i in range(4) + ] + + with _ws_connect(sync_client, sync_client.cookies.get("qm_session")) as ws: + _read_until(ws, "session") + for anomaly in anomalies: + ws.send_json(anomaly) + first = _read_until(ws, "entity")["entity"]["name"] + + with _ws_connect(sync_client, sync_client.cookies.get("qm_session")) as ws: + _read_until(ws, "session") + for anomaly in anomalies: + ws.send_json(anomaly) + second_frame = _read_until(ws, "entity") + assert second_frame["entity"]["name"] == first + assert second_frame["is_new"] is False + assert second_frame["entity"]["contact_count"] == 2