From 9991a836fe7375036b4347fc6ce0f7942f2cf3e1 Mon Sep 17 00:00:00 2001 From: drjones Date: Tue, 6 Oct 2026 23:43:30 -0700 Subject: [PATCH] Snapshot: full project state --- .gitignore | 17 + README.md | 66 ++++ app.py | 587 ++++++++++++++++++++++++++++ collector-stt.service | 17 + collector-tunnel.service | 13 + collector.service | 15 + config.yaml | 52 +++ install.sh | 17 + multiturn_test.py | 62 +++ nightmare_stt.py | 41 ++ selftest.py | 24 ++ static/index.html | 325 +++++++++++++++ static/index.html.bak-bmac-20260923 | 324 +++++++++++++++ tunnel.sh | 30 ++ voice_test.py | 52 +++ 15 files changed, 1642 insertions(+) create mode 100644 .gitignore create mode 100644 README.md create mode 100644 app.py create mode 100644 collector-stt.service create mode 100644 collector-tunnel.service create mode 100644 collector.service create mode 100644 config.yaml create mode 100644 install.sh create mode 100644 multiturn_test.py create mode 100644 nightmare_stt.py create mode 100644 selftest.py create mode 100644 static/index.html create mode 100644 static/index.html.bak-bmac-20260923 create mode 100644 tunnel.sh create mode 100644 voice_test.py diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..24b2937 --- /dev/null +++ b/.gitignore @@ -0,0 +1,17 @@ +__pycache__/ +*.pyc +node_modules/ +.venv/ +venv/ +.env +*.db +*.sqlite* +*.log +.DS_Store +out/ +work/ +.pio/ +briefs/ +dns-backup/ +archive/ +*.png diff --git a/README.md b/README.md new file mode 100644 index 0000000..fc4643c --- /dev/null +++ b/README.md @@ -0,0 +1,66 @@ +# The Collector + +A smooth-operator conversational voice AI running on drjones' homelab. Lives on +Proxmox LXC **CT 771** @ `10.30.20.172:8766`. + +## What it is +A latency-optimized voice assistant. It listens, responds naturally through a +local LLM, speaks with a natural voice, and can make outbound phone calls on +request. + +## Pipeline (latency budget) +| Stage | Engine | Latency | +|-------|--------|---------| +| STT | faster-whisper `small.en` (GPU, nightmare `.29:8770`) | ~0.35s / utterance | +| LLM | `ornith-1.5:9b-64k` @ nightmare `.29:11434` (fallback `.222`) | ~600ms first token, streamed | +| TTS | Piper `lessac-high` (CPU, 22.05k → 16k) | ~400ms / sentence, overlapped w/ LLM | + +Sentence-level TTS/LLM overlap: audio starts flowing after the first sentence +completes while the model keeps generating. Server-side energy VAD with barge-in. + +## Interfaces +- **Browser voice** — `GET /` (WebSocket `/ws`, PCM16 16 kHz full-duplex) +- **Phone** — SignalWire (outbound via `POST /api/call {to, task}`; inbound stays + on JONES's number for now). Media stream `/signalwire/stream` (mulaw 8 kHz). +- **Text** — `POST /api/text {text}` +- **Health** — `GET /api/health` + +## Public URL +Cloudflare quick tunnel (ephemeral): `https://placed-ivory-unto-effective.trycloudflare.com` +Self-registers via `POST /api/config {public_base_url}` (systemd `collector-tunnel.service`). + +## Key decisions +- **LLM lane = nightmare `.29`, not `.222`**: `.222` holds ornith resident but is + saturated by ~92 cron jobs + Kalshi bots and hangs on `/api/chat`; `.29` answers + in ~1s. `keep_alive:30m` keeps ornith warm there. +- **STT on GPU (nightmare), not the CT CPU**: whisper on the CT's 2 cores was 7-8s + (the real latency killer); `faster-whisper` on nightmare's 4080S + (`collector-stt.service` @ `:8770`, CUDA) is ~0.35s — 20×. TTS stays on CPU (Piper + ~400ms, not the bottleneck). +- **Phone** reuses the SignalWire account + number (`+1 254 513 9895`) JONES uses; + outbound TwiML points at this app's tunnel. + +## Files +``` +/opt/collector/ + app.py self-contained FastAPI app (STT/TTS/LLM/phone) + config.yaml model lane, TTS voice, SignalWire creds, persona + selftest.py Piper→whisper round-trip self-test + static/index.html browser voice UI (AudioWorklet) + tunnel.sh cloudflared quick tunnel + self-registration + collector.service systemd unit (app) + collector-tunnel.service systemd unit (tunnel) +``` + +## Ops +```bash +# restart +ssh root@10.30.20.85 "pct exec 771 -- systemctl restart collector" +# logs +ssh root@10.30.20.85 "pct exec 771 -- journalctl -u collector -n 30" +# self-test the audio pipeline +ssh root@10.30.20.85 "pct exec 771 -- python3 /opt/collector/selftest.py" +# make a phone call +curl -s -X POST http://10.30.20.172:8766/api/call \ + -H 'Content-Type: application/json' -d '{"to":"+1XXXXXXXXXX","task":"..."}' +``` diff --git a/app.py b/app.py new file mode 100644 index 0000000..da6e719 --- /dev/null +++ b/app.py @@ -0,0 +1,587 @@ +#!/usr/bin/env python3 +"""The Collector — a smooth-operator conversational voice AI. + +Pipeline (latency-optimized): + STT faster-whisper small.en (CPU, int8) ~1-2s / short utterance + LLM ornith-1.5:9b-64k @ .222 (warm) ~600ms first token, streamed + TTS Piper lessac-high (CPU, 22.05k->16k) ~400ms / sentence, overlapped w/ LLM + +Interfaces: + WS /ws browser full-duplex voice (PCM16 16 kHz) + WS /signalwire/stream phone media stream (mulaw 8 kHz) + POST /api/call {to, task} outbound phone call on request + POST /api/text {text} text chat (returns text) + GET /api/health +""" +import asyncio +import base64 +import json +import logging +import os +import re + +import audioop +import httpx +import numpy as np +import yaml +from fastapi import FastAPI, Request, WebSocket, WebSocketDisconnect +from fastapi.responses import FileResponse, JSONResponse, PlainTextResponse + +logging.basicConfig(level=logging.INFO, + format="%(asctime)s %(name)s %(levelname)s %(message)s") +log = logging.getLogger("collector") + +BASE = "/opt/collector" + + +def load_config() -> dict: + with open(os.path.join(BASE, "config.yaml")) as f: + return yaml.safe_load(f) + + +CFG = load_config() + +# --------------------------------------------------------------------------- +# audio helpers (stdlib audioop) +# --------------------------------------------------------------------------- +def ulaw_to_pcm(data: bytes) -> bytes: + return audioop.ulaw2lin(data, 2) + + +def pcm_to_ulaw(pcm: bytes) -> bytes: + return audioop.lin2ulaw(pcm, 2) + + +def resample(pcm: bytes, src: int, dst: int) -> bytes: + if src == dst or not pcm: + return pcm + out, _ = audioop.ratecv(pcm, 2, 1, src, dst, None) + return out + + +def rms(pcm: bytes) -> int: + return audioop.rms(pcm, 2) if pcm else 0 + + +# --------------------------------------------------------------------------- +# STT +# --------------------------------------------------------------------------- +class STT: + def __init__(self, cfg): + self.cfg = cfg + self.model = None + self.remote_url = (cfg.get("remote_url") or "").rstrip("/") + + def load(self): + # Remote GPU STT (nightmare) is primary; local whisper is a lazy fallback. + return self + + def transcribe(self, pcm16k: bytes) -> str: + if self.remote_url: + try: + return self._remote(pcm16k) + except Exception as e: + log.warning("remote STT failed (%s); falling back to local", e) + return self._local(pcm16k) + + def _remote(self, pcm16k: bytes) -> str: + import requests + r = requests.post(self.remote_url + "/transcribe", data=pcm16k, timeout=30) + r.raise_for_status() + return (r.json().get("text") or "").strip() + + def _local(self, pcm16k: bytes) -> str: + if self.model is None: + from faster_whisper import WhisperModel + c = self.cfg + log.info("Loading local faster-whisper %s (fallback)...", c["model"]) + self.model = WhisperModel(c["model"], device=c["device"], + compute_type=c["compute_type"], + cpu_threads=c.get("cpu_threads", 2)) + audio = np.frombuffer(pcm16k, dtype=np.int16).astype(np.float32) / 32768.0 + segs, _ = self.model.transcribe( + audio, beam_size=1, language="en", vad_filter=True, + condition_on_previous_text=False, no_speech_threshold=0.6) + return " ".join(s.text.strip() for s in segs).strip() + + +# --------------------------------------------------------------------------- +# TTS +# --------------------------------------------------------------------------- +class TTS: + def __init__(self, cfg): + self.cfg = cfg + self.voice_path = cfg.get("piper_voice") or \ + "/opt/collector/voices/en_US-lessac-high.onnx" + self._piper = None + + def load(self): + if self._piper is None: + # onnxruntime pins threads to CPUs 0..N-1; the unprivileged LXC + # cpuset (e.g. 8,17) rejects that. Setting an explicit thread count + # disables the affinity pin (which is otherwise just log noise). + import onnxruntime + _orig = onnxruntime.InferenceSession.__init__ + + def _patched(self, *a, **k): + so = k.get("sess_options") or onnxruntime.SessionOptions() + so.intra_op_num_threads = 2 + k["sess_options"] = so + _orig(self, *a, **k) + + onnxruntime.InferenceSession.__init__ = _patched + + from piper import PiperVoice + log.info("Loading Piper voice %s ...", self.voice_path) + self._piper = PiperVoice.load(self.voice_path) + log.info("TTS ready") + return self._piper + + def synthesize_16k(self, text: str) -> bytes: + """Return PCM16 mono 16 kHz for `text` (piper yields AudioChunk iterators).""" + if not text.strip(): + return b"" + speed = float(self.cfg.get("speed", 1.05)) + from piper import SynthesisConfig + syn = SynthesisConfig(length_scale=1.0 / speed) + chunks = [] + rate = 22050 + for chunk in self._piper.synthesize(text, syn_config=syn): + audio = (chunk.audio_float_array * 32767).astype(np.int16).tobytes() + rate = chunk.sample_rate + chunks.append(audio) + pcm = b"".join(chunks) + return resample(pcm, rate, 16000) + + +# --------------------------------------------------------------------------- +# LLM (ornith streaming, /api/chat, think:false) +# --------------------------------------------------------------------------- +async def stream_llm(messages, cfg): + payload = { + "model": cfg["ollama"]["model"], + "messages": messages, + "stream": True, + "think": False, + "keep_alive": "30m", + "options": { + "temperature": cfg["ollama"].get("temperature", 0.5), + "num_predict": cfg["ollama"].get("num_predict", 256), + }, + } + urls = [cfg["ollama"]["base_url"], cfg["ollama"].get("fallback_url", "")] + for base in urls: + if not base: + continue + url = base.rstrip("/") + "/api/chat" + try: + async with httpx.AsyncClient( + timeout=httpx.Timeout(10.0, read=60.0)) as client: + async with client.stream("POST", url, json=payload) as r: + r.raise_for_status() + async for line in r.aiter_lines(): + if not line.strip(): + continue + try: + obj = json.loads(line) + except Exception: + continue + if obj.get("done"): + return + msg = obj.get("message", {}) + content = msg.get("content") or msg.get("thinking") or "" + if content: + yield content + return + except Exception as e: + log.warning("LLM host %s failed: %s", base, e) + log.error("All LLM hosts failed") + + +async def sentence_stream(chunks): + """Yield sentences as they complete, enabling TTS/LLM overlap.""" + buf = "" + async for chunk in chunks: + buf += chunk + while True: + m = re.search(r"[.!?\n]", buf) + if not m: + break + idx = m.end() + sentence = buf[:idx].strip() + buf = buf[idx:].strip() + if sentence: + yield sentence + if buf.strip(): + yield buf.strip() + + +# --------------------------------------------------------------------------- +# Conversation session (shared by browser WS + phone stream) +# --------------------------------------------------------------------------- +SPEECH_RMS = 400 # 16-bit PCM: voice onset +BARGE_RMS = 900 # sustained level that interrupts playback +SILENCE_END_MS = 700 # ms of quiet after speech before we turn +FRAME_MS = 20 +MAX_UTT_BYTES = 18 * 16000 * 2 # 18s cap + + +class Session: + def __init__(self, cfg, system_prompt): + self.cfg = cfg + self.stt = STT(cfg["stt"]) + self.tts = TTS(cfg["tts"]) + self.messages = [{"role": "system", "content": system_prompt}] + self.buf = bytearray() + self.talking = False + self.silence_ms = 0 + self._barge = False + self._processing = False + self._queue = asyncio.Queue() + self._worker = None + + async def warm(self): + await asyncio.to_thread(self.stt.load) + await asyncio.to_thread(self.tts.load) + try: + async for _ in stream_llm( + [{"role": "user", "content": "hi"}], self.cfg): + break + except Exception as e: + log.warning("LLM warmup: %s", e) + log.info("Collector warmed (STT + TTS + LLM)") + + def start_worker(self, send_audio, send_text): + """Begin the continuous conversation loop for this connection.""" + self._worker = asyncio.create_task( + self._process_loop(send_audio, send_text)) + + async def feed_audio(self, pcm16k: bytes): + """Continuous VAD. Call for every 20ms frame. Never drops speech — + even during a reply, overlapping audio is buffered and processed next + (barge-in).""" + level = rms(pcm16k) + if level > SPEECH_RMS: + self.talking = True + self.silence_ms = 0 + if self._processing: + self._barge = True # user started talking over the reply + elif self.talking: + self.silence_ms += FRAME_MS + self.buf.extend(pcm16k) + if self.talking and (self.silence_ms >= SILENCE_END_MS or + len(self.buf) > MAX_UTT_BYTES): + self.talking = False + self.silence_ms = 0 + pcm = bytes(self.buf) + self.buf.clear() + self._queue.put_nowait(pcm) + + async def _process_loop(self, send_audio, send_text): + while True: + pcm = await self._queue.get() + self._processing = True + self._barge = False + text = await asyncio.to_thread(self.stt.transcribe, pcm) + if text: + await send_text(text, "user") + self.messages.append({"role": "user", "content": text}) + await self._reply(send_audio, send_text) + self._processing = False + + async def speak(self, text: str, send_audio, send_text): + """Speak an injected message (outbound preamble), then keep listening.""" + self.messages.append({"role": "user", "content": text}) + self._processing = True + self._barge = False + await self._reply(send_audio, send_text) + self._processing = False + + async def _reply(self, send_audio, send_text): + full = [] + async for sentence in sentence_stream(stream_llm(self.messages, self.cfg)): + if self._barge: + log.info("barge-in: stopping reply") + break + full.append(sentence) + pcm = await asyncio.to_thread(self.tts.synthesize_16k, sentence) + if pcm: + await send_audio(pcm) + if full: + reply = " ".join(full) + await send_text(reply, "assistant") + self.messages.append({"role": "assistant", "content": reply}) + if len(self.messages) > 24: + self.messages = [self.messages[0]] + self.messages[-23:] + + +# --------------------------------------------------------------------------- +# SignalWire client +# --------------------------------------------------------------------------- +class Phone: + def __init__(self, cfg): + t = cfg.get("signalwire", {}) + self.sid = t.get("account_sid", "") + self.token = t.get("auth_token", "") + self.number = t.get("phone_number", "") + self.client = None + self.tasks = {} # call_sid -> outbound task + if self.sid and self.token: + from signalwire.rest import Client as SWClient + self.client = SWClient(self.sid, self.token, + signalwire_space_url="templeofdoom.signalwire.com") + + @property + def ready(self): + return self.client is not None and bool(self.number) + + def create_call(self, to: str, url: str): + return self.client.calls.create(to=to, from_=self.number, url=url) + + +PHONE = Phone(CFG) +app = FastAPI(title="The Collector") + + +def system_prompt() -> str: + return CFG["agent"].get("system_prompt", "You are The Collector.") + + +def stream_twiml(direction: str, public_base: str) -> str: + # SignalWire requires wss://, not https:// + wss_base = public_base.replace("https://", "wss://", 1) + stream_url = f"{wss_base}/signalwire/stream?direction={direction}" + return (f'' + f'') + + +# --------------------------------------------------------------------------- +# Health +# --------------------------------------------------------------------------- +@app.get("/") +async def index(): + return FileResponse(os.path.join(BASE, "static", "index.html")) + + +@app.get("/api/health") +async def health(): + return { + "ok": True, + "name": CFG["agent"]["name"], + "model": CFG["ollama"]["model"], + "llm": CFG["ollama"]["base_url"], + "phone": PHONE.ready, + "number": PHONE.number if PHONE.ready else None, + } + + +@app.post("/api/config") +async def api_config(request: Request): + """Accept runtime config (tunnel URL, temperature, tts speed) + persist.""" + body = await request.json() + changed = {} + if "public_base_url" in body: + url = (body["public_base_url"] or "").strip().rstrip("/") + CFG.setdefault("server", {})["public_base_url"] = url + changed["public_base_url"] = url + if "temperature" in body: + t = float(body["temperature"]) + CFG.setdefault("ollama", {})["temperature"] = max(0.0, min(1.5, t)) + changed["temperature"] = CFG["ollama"]["temperature"] + if "speed" in body: + s = float(body["speed"]) + CFG.setdefault("tts", {})["speed"] = max(0.5, min(2.0, s)) + changed["speed"] = CFG["tts"]["speed"] + if changed: + try: + with open(os.path.join(BASE, "config.yaml")) as f: + cfg = yaml.safe_load(f) + if "public_base_url" in changed: + cfg.setdefault("server", {})["public_base_url"] = changed["public_base_url"] + if "temperature" in changed: + cfg.setdefault("ollama", {})["temperature"] = changed["temperature"] + if "speed" in changed: + cfg.setdefault("tts", {})["speed"] = changed["speed"] + with open(os.path.join(BASE, "config.yaml"), "w") as f: + yaml.safe_dump(cfg, f, sort_keys=False) + except Exception as e: + log.warning("config persist failed: %s", e) + return {"ok": True, **changed} + return JSONResponse({"ok": False, "error": "nothing to set"}, 400) + + +@app.get("/api/settings") +async def api_settings(): + return { + "temperature": CFG["ollama"].get("temperature", 0.5), + "speed": CFG["tts"].get("speed", 1.05), + "model": CFG["ollama"]["model"], + "voice": "lessac-high", + "llm": CFG["ollama"]["base_url"], + "stt": "GPU (nightmare)", + } + + +# --------------------------------------------------------------------------- +# Text chat +# --------------------------------------------------------------------------- +@app.post("/api/text") +async def api_text(request: Request): + body = await request.json() + text = (body.get("text") or "").strip() + if not text: + return JSONResponse({"ok": False, "error": "empty text"}, 400) + msgs = [{"role": "system", "content": system_prompt()}, + {"role": "user", "content": text}] + out = [] + async for chunk in stream_llm(msgs, CFG): + out.append(chunk) + return {"ok": True, "reply": "".join(out).strip()} + + +# --------------------------------------------------------------------------- +# Outbound call +# --------------------------------------------------------------------------- +@app.post("/api/call") +async def api_call(request: Request): + if not PHONE.ready: + return JSONResponse({"ok": False, "error": "SignalWire not configured"}, 400) + body = await request.json() + to = (body.get("to") or "").strip() + task = (body.get("task") or "").strip() + if not to: + return JSONResponse({"ok": False, "error": "missing 'to'"}, 400) + public = CFG["server"].get("public_base_url", "") + if not public: + return JSONResponse({"ok": False, "error": "public_base_url not set"}, 500) + call = await asyncio.to_thread( + PHONE.create_call, to, f"{public}/signalwire/outbound") + PHONE.tasks[call.sid] = task + return {"ok": True, "call_sid": call.sid, "to": to, "task": task} + + +# --------------------------------------------------------------------------- +# SignalWire webhooks +# --------------------------------------------------------------------------- +@app.post("/signalwire/voice") +async def sw_inbound(): + public = CFG["server"].get("public_base_url", "") + return PlainTextResponse(stream_twiml("inbound", public), + media_type="application/xml") + + +@app.post("/signalwire/outbound") +async def sw_outbound(): + public = CFG["server"].get("public_base_url", "") + return PlainTextResponse(stream_twiml("outbound", public), + media_type="application/xml") + + +# --------------------------------------------------------------------------- +# Browser full-duplex voice +# --------------------------------------------------------------------------- +@app.websocket("/ws") +async def browser_ws(ws: WebSocket): + await ws.accept() + session = Session(CFG, system_prompt()) + + async def send_audio(pcm): + await ws.send_bytes(pcm) + + async def send_text(text, role): + await ws.send_text(json.dumps({"type": "transcript", "role": role, + "text": text})) + + await session.warm() + session.start_worker(send_audio, send_text) + try: + while True: + data = await ws.receive() + if data["type"] == "websocket.receive": + if data.get("bytes"): + await session.feed_audio(data["bytes"]) + elif data["type"] == "websocket.disconnect": + break + except WebSocketDisconnect: + pass + finally: + if session._worker: + session._worker.cancel() + + +# --------------------------------------------------------------------------- +# SignalWire media stream (phone) +# --------------------------------------------------------------------------- +async def send_phone_audio(ws, stream_sid, pcm16k): + pcm8 = resample(pcm16k, 16000, 8000) + ulaw = pcm_to_ulaw(pcm8) + frame = 160 # 20ms @ 8kHz + for i in range(0, len(ulaw), frame): + chunk = ulaw[i:i + frame] + await ws.send_text(json.dumps({ + "event": "media", "streamSid": stream_sid, + "media": {"payload": base64.b64encode(chunk).decode()}})) + await asyncio.sleep(0.02) + + +@app.websocket("/signalwire/stream") +async def sw_stream(ws: WebSocket, direction: str = "inbound"): + await ws.accept() + session = Session(CFG, system_prompt()) + stream_sid = None + call_sid = None + first = True + + async def send_audio(pcm16k): + await send_phone_audio(ws, stream_sid, pcm16k) + + async def send_text(text, role): + pass # no transcript channel on phone; audio only + + await session.warm() + session.start_worker(send_audio, send_text) + try: + while True: + msg = await ws.receive_text() + try: + obj = json.loads(msg) + except Exception: + continue + ev = obj.get("event") + + if ev == "start": + start = obj.get("start", {}) + stream_sid = start.get("streamSid", stream_sid) + call_sid = start.get("callSid", call_sid) + if direction == "outbound" and first: + first = False + task = PHONE.tasks.pop(call_sid, "") + if task: + preamble = CFG["agent"].get("outbound_preamble", "") + task + await session.speak(preamble, send_audio, send_text) + + elif ev == "media": + payload = obj.get("media", {}).get("payload", "") + if not payload: + continue + raw = base64.b64decode(payload) + pcm8 = ulaw_to_pcm(raw) + pcm16k = resample(pcm8, 8000, 16000) + await session.feed_audio(pcm16k) + + elif ev == "stop": + break + except WebSocketDisconnect: + pass + finally: + if session._worker: + session._worker.cancel() + + +# --------------------------------------------------------------------------- +# Entrypoint +# --------------------------------------------------------------------------- +if __name__ == "__main__": + import uvicorn + s = CFG["server"] + uvicorn.run(app, host=s.get("host", "0.0.0.0"), port=s.get("port", 8766)) diff --git a/collector-stt.service b/collector-stt.service new file mode 100644 index 0000000..4c7ed6e --- /dev/null +++ b/collector-stt.service @@ -0,0 +1,17 @@ +[Unit] +Description=The Collector — GPU STT (faster-whisper on CUDA) +After=network-online.target + +[Service] +Type=simple +User=drjones +Group=drjones +ExecStart=/usr/bin/python3 /opt/collector-stt/nightmare_stt.py +WorkingDirectory=/opt/collector-stt +Restart=always +RestartSec=5 +Environment=PYTHONUNBUFFERED=1 +Environment=LD_LIBRARY_PATH=/home/drjones/.local/lib/python3.14/site-packages/nvidia/cuda_runtime/lib:/home/drjones/.local/lib/python3.14/site-packages/nvidia/cudnn/lib:/home/drjones/.local/lib/python3.14/site-packages/nvidia/cublas/lib:/home/drjones/.local/lib/python3.14/site-packages/nvidia/cuda_nvrtc/lib + +[Install] +WantedBy=multi-user.target diff --git a/collector-tunnel.service b/collector-tunnel.service new file mode 100644 index 0000000..8e473a4 --- /dev/null +++ b/collector-tunnel.service @@ -0,0 +1,13 @@ +[Unit] +Description=The Collector — Cloudflare quick tunnel (public URL + phone webhooks) +After=collector.service network-online.target +Wants=network-online.target + +[Service] +Type=simple +ExecStart=/bin/bash /opt/collector/tunnel.sh +Restart=always +RestartSec=10 + +[Install] +WantedBy=multi-user.target diff --git a/collector.service b/collector.service new file mode 100644 index 0000000..e5f7a9c --- /dev/null +++ b/collector.service @@ -0,0 +1,15 @@ +[Unit] +Description=The Collector — conversational voice AI (ornith + Piper + SignalWire) +After=network-online.target +Wants=network-online.target + +[Service] +Type=simple +ExecStart=/usr/bin/python3 /opt/collector/app.py +WorkingDirectory=/opt/collector +Restart=always +RestartSec=5 +Environment=PYTHONUNBUFFERED=1 + +[Install] +WantedBy=multi-user.target diff --git a/config.yaml b/config.yaml new file mode 100644 index 0000000..310c53a --- /dev/null +++ b/config.yaml @@ -0,0 +1,52 @@ +server: + host: 0.0.0.0 + port: 8766 + public_base_url: "" + +# LLM — ornith conversational lane. +# Primary = nightmare .29 (4080S, measured fast first-token; .222 is saturated by +# the 92 cron jobs + Kalshi bots and hangs on /api/chat). Fallback = .222 (ornith home lane). +ollama: + base_url: http://10.30.20.29:11434 + fallback_url: http://10.30.20.222:11434 + model: ornith-1.5:9b-64k + think: false + temperature: 0.5 + num_predict: 256 + +stt: + model: small.en + device: cpu + compute_type: int8 + cpu_threads: 2 + remote_url: http://10.30.20.29:8770 + +# Piper lessac-high = natural voice, ~400ms/sentence on CPU, no GPU needed. +# (TTS is NOT the latency bottleneck — STT + LLM first-token dominate the budget.) +tts: + engine: piper + voice: lessac-high + speed: 1.05 + piper_voice: /opt/collector/voices/en_US-lessac-high.onnx + +agent: + name: The Collector + owner_phone: "+14252800023" + system_prompt: > + You are The Collector, a smooth, effortlessly charming conversational AI running on + Indiana's homelab. You are on a LIVE VOICE call. Reply in 1-3 short conversational + sentences. No markdown, no bullet points, no emoji. Listen first, then respond with + warmth and wit — a trusted friend who happens to be exceptionally smooth. Match the + caller's energy. Keep it natural, human, and brief. + outbound_preamble: > + You are The Collector, making a phone call on Indiana's behalf. Introduce yourself + naturally, accomplish the objective, confirm details, wrap up politely. Objective: + +signalwire: + account_sid: 8a89a49f-6e98-4de1-b32e-df897f656f24 + auth_token: PT2ad398d81d9a1632375b8ac88844ea8ecfa46021b5389dfb + phone_number: "+12545139895" + +api: + auth_key: collector-voice-2026 + auth_enabled: false diff --git a/install.sh b/install.sh new file mode 100644 index 0000000..1f2843c --- /dev/null +++ b/install.sh @@ -0,0 +1,17 @@ +#!/bin/bash +# The Collector — dependency install (Debian 12 LXC) +set -e +export DEBIAN_FRONTEND=noninteractive + +apt-get update -qq +apt-get install -y -qq \ + python3 python3-pip python3-venv python3-yaml \ + ffmpeg libsndfile1 espeak-ng + +pip3 install --break-system-packages --no-cache-dir \ + fastapi "uvicorn[standard]" websockets httpx \ + twilio signalwire pyyaml numpy \ + faster-whisper piper-tts requests python-multipart + +mkdir -p /opt/collector/voices +echo "INSTALL_OK" diff --git a/multiturn_test.py b/multiturn_test.py new file mode 100644 index 0000000..17a973c --- /dev/null +++ b/multiturn_test.py @@ -0,0 +1,62 @@ +#!/usr/bin/env python3 +"""Multi-turn conversational test: two back-and-forth exchanges over /ws.""" +import asyncio +import json +import sys + +sys.path.insert(0, "/opt/collector") +from app import TTS, CFG +import websockets + +tts = TTS(CFG["tts"]) +tts.load() +P1 = tts.synthesize_16k("What is your name?") +P2 = tts.synthesize_16k("Tell me something interesting.") + +FRAME = 640 # 20ms @ 16kHz + + +async def send_utterance(ws, pcm): + for i in range(0, len(pcm), FRAME): + await ws.send(pcm[i:i + FRAME]) + await asyncio.sleep(0.02) + silence = b"\x00" * FRAME + for _ in range(50): # 1s silence -> VAD turn end + await ws.send(silence) + await asyncio.sleep(0.02) + + +async def read_reply(ws, seen_assistant): + """Read messages until we collect a NEW assistant transcript + audio.""" + audio = 0 + while True: + try: + msg = await asyncio.wait_for(ws.recv(), timeout=45) + except Exception: + return False + if isinstance(msg, str): + m = json.loads(msg) + if m.get("role") == "assistant": + seen_assistant += 1 + else: + audio += len(msg) + if seen_assistant >= 1 and audio > 0: + return True + + +async def main(): + async with websockets.connect("ws://127.0.0.1:8766/ws") as ws: + await asyncio.sleep(6) # warm + # Turn 1 + await send_utterance(ws, P1) + ok1 = await read_reply(ws, 0) + print("TURN1 reply:", "PASS" if ok1 else "FAIL") + # Turn 2 + await asyncio.sleep(1) + await send_utterance(ws, P2) + ok2 = await read_reply(ws, 1) + print("TURN2 reply:", "PASS" if ok2 else "FAIL") + print("MULTITURN", "PASS" if (ok1 and ok2) else "FAIL") + + +asyncio.run(main()) diff --git a/nightmare_stt.py b/nightmare_stt.py new file mode 100644 index 0000000..8c1075c --- /dev/null +++ b/nightmare_stt.py @@ -0,0 +1,41 @@ +#!/usr/bin/env python3 +"""GPU STT service for The Collector — faster-whisper small.en on CUDA (4080S). + +POST /transcribe (raw PCM16 mono 16 kHz body) -> {"text": "..."} +GET /health +""" +import numpy as np +from fastapi import FastAPI, Request +from faster_whisper import WhisperModel +import uvicorn + +MODEL = "small.en" +app = FastAPI() +model = None + + +@app.on_event("startup") +async def load(): + global model + model = WhisperModel(MODEL, device="cuda", compute_type="float16") + print(f"STT ready: {MODEL} on cuda") + + +@app.get("/health") +async def health(): + return {"ok": model is not None, "model": MODEL, "device": "cuda"} + + +@app.post("/transcribe") +async def transcribe(request: Request): + body = await request.body() + audio = np.frombuffer(body, dtype=np.int16).astype(np.float32) / 32768.0 + segs, _ = model.transcribe( + audio, beam_size=1, language="en", vad_filter=True, + condition_on_previous_text=False, no_speech_threshold=0.6) + text = " ".join(s.text.strip() for s in segs).strip() + return {"text": text} + + +if __name__ == "__main__": + uvicorn.run(app, host="0.0.0.0", port=8770) diff --git a/selftest.py b/selftest.py new file mode 100644 index 0000000..da1072b --- /dev/null +++ b/selftest.py @@ -0,0 +1,24 @@ +#!/usr/bin/env python3 +"""Collector self-test: synthesize a phrase with Piper, transcribe it back with +faster-whisper. Proves the STT/TTS audio pipeline end-to-end (no phone needed).""" +import sys +sys.path.insert(0, "/opt/collector") +from app import STT, TTS, CFG + +stt = STT(CFG["stt"]) +tts = TTS(CFG["tts"]) + +print("loading STT (first run downloads whisper model)...") +stt.load() +print("loading TTS...") +tts.load() + +phrase = "This is the collector self test. If you can hear this, the voice pipeline works." +pcm16k = tts.synthesize_16k(phrase) +print("TTS_OK bytes=%d" % len(pcm16k)) + +text = stt.transcribe(pcm16k) +print("STT_OK transcript=%r" % text) + +ok = "collector" in text.lower() or "voice" in text.lower() or "hear" in text.lower() +print("SELFTEST", "PASS" if ok else "PARTIAL", file=sys.stderr) diff --git a/static/index.html b/static/index.html new file mode 100644 index 0000000..f4c05e7 --- /dev/null +++ b/static/index.html @@ -0,0 +1,325 @@ + + + + + +The Collector + + + +
+ +
+

The Collector

+
smooth operator · local voice
+ +
+ +
+
+ +
offline
+ + + +
+ + +
+ +
+
0.50
+
1.05
+
+ modelornith-1.5:9b-64k +
+
+ voicelessac-high · STT on GPU +
+
+ +
+
+
+ +
+
+ +
+
tap Start and talk — it listens, thinks, and speaks back
+
+ + +☕ Support on Buy Me a Coffee + + diff --git a/static/index.html.bak-bmac-20260923 b/static/index.html.bak-bmac-20260923 new file mode 100644 index 0000000..093159c --- /dev/null +++ b/static/index.html.bak-bmac-20260923 @@ -0,0 +1,324 @@ + + + + + +The Collector + + + +
+ +
+

The Collector

+
smooth operator · local voice
+ +
+ +
+
+ +
offline
+ + + +
+ + +
+ +
+
0.50
+
1.05
+
+ modelornith-1.5:9b-64k +
+
+ voicelessac-high · STT on GPU +
+
+ +
+
+
+ +
+
+ +
+
tap Start and talk — it listens, thinks, and speaks back
+
+ + + + diff --git a/tunnel.sh b/tunnel.sh new file mode 100644 index 0000000..45af300 --- /dev/null +++ b/tunnel.sh @@ -0,0 +1,30 @@ +#!/bin/bash +# Start a Cloudflare quick tunnel for The Collector, extract the trycloudflare +# URL, and self-register it with the app (/api/config). Blocks while the tunnel +# runs (intended for a systemd service). +LOG=/opt/collector/tunnel.log +APP=http://127.0.0.1:8766 + +: > "$LOG" +/usr/local/bin/cloudflared tunnel --url http://localhost:8766 --no-autoupdate >> "$LOG" 2>&1 & +CF_PID=$! + +trap 'kill $CF_PID 2>/dev/null' EXIT + +URL="" +for _ in $(seq 1 45); do + URL=$(grep -oE 'https://[a-z0-9-]+\.trycloudflare\.com' "$LOG" | head -1) + [ -n "$URL" ] && break + sleep 2 +done + +if [ -n "$URL" ]; then + echo "TUNNEL_URL=$URL" + curl -s -X POST "$APP/api/config" -H 'Content-Type: application/json' \ + -d "{\"public_base_url\":\"$URL\"}" >/dev/null 2>&1 + echo "registered with app" +else + echo "TUNNEL_FAILED" +fi + +wait $CF_PID diff --git a/voice_test.py b/voice_test.py new file mode 100644 index 0000000..63a4b49 --- /dev/null +++ b/voice_test.py @@ -0,0 +1,52 @@ +#!/usr/bin/env python3 +"""End-to-end voice test: synthesize speech, push it through the /ws browser +channel, assert we get back a transcript + reply audio (PCM16).""" +import asyncio +import json +import sys + +sys.path.insert(0, "/opt/collector") +from app import TTS, CFG + +import websockets + +tts = TTS(CFG["tts"]) +tts.load() +pcm = tts.synthesize_16k("Hello Collector, tell me who you are in one sentence.") + +FRAME = 640 # 20ms @ 16kHz + + +async def main(): + async with websockets.connect("ws://127.0.0.1:8766/ws") as ws: + await asyncio.sleep(6) # let server warm (STT+TTS+LLM) + # speech in 20ms frames + for i in range(0, len(pcm), FRAME): + await ws.send(pcm[i:i + FRAME]) + await asyncio.sleep(0.02) + # 1s silence -> trigger VAD turn end + silence = b"\x00" * FRAME + for _ in range(50): + await ws.send(silence) + await asyncio.sleep(0.02) + + got_text = got_audio = False + deadline = asyncio.get_event_loop().time() + 60 + while asyncio.get_event_loop().time() < deadline: + try: + msg = await asyncio.wait_for(ws.recv(), timeout=40) + except Exception: + break + if isinstance(msg, str): + print("TRANSCRIPT:", msg) + got_text = True + else: + print("AUDIO bytes:", len(msg)) + got_audio = True + if got_text and got_audio: + break + + print("VOICE_TEST", "PASS" if (got_text and got_audio) else "FAIL") + + +asyncio.run(main())