feat: pipeline FastAPI app — WebSocket audio ingestion + internal mute API
This commit is contained in:
@@ -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)
|
||||
Reference in New Issue
Block a user