From 9a672c4c17e003a9c1640f185a951f559bf3379d Mon Sep 17 00:00:00 2001 From: Pluto Date: Thu, 28 May 2026 10:39:04 -0500 Subject: [PATCH] =?UTF-8?q?feat:=20pipeline=20FastAPI=20app=20=E2=80=94=20?= =?UTF-8?q?WebSocket=20audio=20ingestion=20+=20internal=20mute=20API?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/pipeline/main.py | 101 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 101 insertions(+) create mode 100644 app/pipeline/main.py diff --git a/app/pipeline/main.py b/app/pipeline/main.py new file mode 100644 index 0000000..ae1a62c --- /dev/null +++ b/app/pipeline/main.py @@ -0,0 +1,101 @@ +import asyncio +import json +from contextlib import asynccontextmanager +from datetime import datetime, timezone + +import uvicorn +from fastapi import FastAPI, WebSocket, WebSocketDisconnect + +from app.shared.config import get_settings, Settings +from app.shared.database import get_db, init_schema +from app.pipeline.ingestion import get_or_create_room, mute_room, unmute_room, is_muted +from app.pipeline.processor import process_utterance + +_settings: Settings | None = None + + +def _get_settings() -> Settings: + global _settings + if _settings is None: + _settings = get_settings() + return _settings + + +@asynccontextmanager +async def lifespan(app: FastAPI): + s = _get_settings() + init_schema(s) + yield + + +app = FastAPI(lifespan=lifespan) + + +@app.websocket("/ws/{room_name}") +async def audio_ws(websocket: WebSocket, room_name: str): + await websocket.accept() + s = _get_settings() + db = get_db(s) + + # Upsert room in DB + db.execute( + "INSERT INTO rooms (name, is_active, last_seen) VALUES (?,1,CURRENT_TIMESTAMP) " + "ON CONFLICT(name) DO UPDATE SET is_active=1, last_seen=CURRENT_TIMESTAMP", + (room_name,), + ) + db.commit() + room = db.execute("SELECT id FROM rooms WHERE name=?", (room_name,)).fetchone() + room_id = room["id"] + + get_or_create_room(room_name, s.vad_silence_ms, s.vad_min_speech_ms) + + try: + while True: + data = await websocket.receive() + if "bytes" in data: + chunk = data["bytes"] + if not is_muted(room_name): + db.execute( + "UPDATE rooms SET last_seen=CURRENT_TIMESTAMP WHERE id=?", + (room_id,), + ) + db.commit() + now = datetime.now(timezone.utc) + room_state = get_or_create_room(room_name, s.vad_silence_ms, s.vad_min_speech_ms) + result = room_state["buffer"].process_chunk(chunk, now) + if result is not None: + audio_bytes, start_time = result + end_time = datetime.now(timezone.utc) + asyncio.create_task( + process_utterance(room_id, audio_bytes, start_time, end_time, s, db) + ) + elif "text" in data: + msg = json.loads(data["text"]) + if msg.get("type") == "mute": + mute_room(room_name) + elif msg.get("type") == "unmute": + unmute_room(room_name) + except WebSocketDisconnect: + db.execute("UPDATE rooms SET is_active=0 WHERE id=?", (room_id,)) + db.commit() + + +@app.post("/internal/rooms/{room_name}/mute") +async def internal_mute(room_name: str): + mute_room(room_name) + return {"status": "muted", "room": room_name} + + +@app.post("/internal/rooms/{room_name}/unmute") +async def internal_unmute(room_name: str): + unmute_room(room_name) + return {"status": "unmuted", "room": room_name} + + +@app.get("/health") +async def health(): + return {"status": "ok"} + + +if __name__ == "__main__": + uvicorn.run("app.pipeline.main:app", host="0.0.0.0", port=8300, reload=False)