- LiteLLM proxy with langfuse callbacks, postgres backends, and OpenRouter fallbacks. - Langfuse observability UI. - Pipecat voice pipeline (LiveKit + STT + TTS + LLM) and Silero TTS build contexts. - Ollama tuned for GPU (OLLAMA_NUM_GPU=999, mem_limit=4g, max 2 loaded models). - open-webui wired to litellm + faster-whisper + silero for voice. - litellm-config.yaml publishes oO's model aliases (tip-generator, embedder, judge) pointing at the host ollama on :11434 so ml/serving can call them via LiteLLM. .env skipped (secrets). Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
229 lines
8.9 KiB
Python
229 lines
8.9 KiB
Python
import asyncio
|
||
import os
|
||
import re
|
||
import uuid
|
||
import logging
|
||
|
||
from fastapi import FastAPI
|
||
from fastapi.responses import HTMLResponse
|
||
from fastapi.staticfiles import StaticFiles
|
||
from pydantic import BaseModel
|
||
|
||
from livekit import api as lkapi
|
||
|
||
from pipecat.audio.vad.silero import SileroVADAnalyzer
|
||
from pipecat.audio.vad.vad_analyzer import VADParams
|
||
from pipecat.frames.frames import TextFrame
|
||
from pipecat.pipeline.pipeline import Pipeline
|
||
from pipecat.pipeline.runner import PipelineRunner
|
||
from pipecat.pipeline.task import PipelineParams, PipelineTask
|
||
from pipecat.processors.aggregators.openai_llm_context import OpenAILLMContext
|
||
from pipecat.processors.frame_processor import FrameProcessor, FrameDirection
|
||
from pipecat.services.openai.llm import OpenAILLMService
|
||
from pipecat.services.openai.stt import OpenAISTTService
|
||
from pipecat.services.openai.tts import OpenAITTSService
|
||
from pipecat.transports.livekit.transport import LiveKitTransport, LiveKitParams
|
||
|
||
|
||
# ── TTS text normalizer ──────────────────────────────────────────────────────
|
||
# Replaces symbols and abbreviations with spoken Russian words so Silero TTS
|
||
# doesn't truncate on unknown characters.
|
||
|
||
_NORM_RULES: list[tuple[re.Pattern, str]] = [
|
||
# Temperature: +12°C / -5°С / 12 °C → плюс двенадцать градусов цельсия
|
||
(re.compile(r"([+-]?\d+)\s*°\s*[CСcс]", re.IGNORECASE), r"\1 градусов цельсия"),
|
||
# Bare degree sign: 90° → 90 градусов
|
||
(re.compile(r"(\d+)\s*°"), r"\1 градусов"),
|
||
# Percent
|
||
(re.compile(r"(\d+)\s*%"), r"\1 процентов"),
|
||
# Speed: m/s, м/с, km/h, км/ч
|
||
(re.compile(r"\bm/s\b", re.IGNORECASE), "метров в секунду"),
|
||
(re.compile(r"\bм/с\b"), "метров в секунду"),
|
||
(re.compile(r"\bkm/h\b", re.IGNORECASE), "километров в час"),
|
||
(re.compile(r"\bкм/ч\b"), "километров в час"),
|
||
# Currency
|
||
(re.compile(r"\$\s*(\d+)"), r"\1 долларов"),
|
||
(re.compile(r"(\d+)\s*\$"), r"\1 долларов"),
|
||
(re.compile(r"€\s*(\d+)"), r"\1 евро"),
|
||
(re.compile(r"(\d+)\s*€"), r"\1 евро"),
|
||
(re.compile(r"(\d+)\s*₽"), r"\1 рублей"),
|
||
# Plus/minus signs before numbers
|
||
(re.compile(r"\+(\d)"), r"плюс \1"),
|
||
(re.compile(r"-(\d)"), r"минус \1"),
|
||
# Common abbreviations
|
||
(re.compile(r"\bкг\b"), "килограмм"),
|
||
(re.compile(r"\bг\b(?=\s|$)"), "грамм"),
|
||
(re.compile(r"\bмм\b"), "миллиметров"),
|
||
(re.compile(r"\bсм\b"), "сантиметров"),
|
||
(re.compile(r"\bкм\b"), "километров"),
|
||
# Strip remaining special chars that TTS can't handle
|
||
(re.compile(r"[°•·†‡§¶©®™«»<>{}[\]|\\~^`]"), ""),
|
||
]
|
||
|
||
|
||
def normalize_for_tts(text: str) -> str:
|
||
"""Replace symbols with spoken Russian equivalents."""
|
||
for pattern, replacement in _NORM_RULES:
|
||
text = pattern.sub(replacement, text)
|
||
return text
|
||
|
||
|
||
class TTSTextNormalizer(FrameProcessor):
|
||
"""Intercepts TextFrames between LLM and TTS, normalizing symbols to words."""
|
||
|
||
async def process_frame(self, frame, direction: FrameDirection = FrameDirection.DOWNSTREAM):
|
||
await super().process_frame(frame, direction)
|
||
if isinstance(frame, TextFrame):
|
||
original = frame.text
|
||
normalized = normalize_for_tts(original)
|
||
if normalized != original:
|
||
logger.debug(f"TTSTextNormalizer: {original!r} → {normalized!r}")
|
||
await self.push_frame(TextFrame(text=normalized), direction)
|
||
else:
|
||
await self.push_frame(frame, direction)
|
||
|
||
logging.basicConfig(level=logging.INFO)
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# ── Config ────────────────────────────────────────────────────────────────────
|
||
LIVEKIT_URL = os.getenv("LIVEKIT_URL", "ws://host.docker.internal:7880") # bot connects here
|
||
LIVEKIT_PUBLIC_URL = os.getenv("LIVEKIT_PUBLIC_URL", "wss://lk.alogins.net") # browser connects here
|
||
LIVEKIT_API_KEY = os.getenv("LIVEKIT_API_KEY", "devkey")
|
||
LIVEKIT_SECRET = os.getenv("LIVEKIT_SECRET", "")
|
||
ADOLF_URL = os.getenv("ADOLF_URL", "http://host.docker.internal:8000/v1")
|
||
STT_URL = os.getenv("STT_URL", "http://host.docker.internal:8880/v1")
|
||
TTS_URL = os.getenv("TTS_URL", "http://host.docker.internal:8881/v1")
|
||
STT_MODEL = os.getenv("STT_MODEL", "deepdml/faster-whisper-large-v3-turbo-ct2")
|
||
TTS_VOICE = os.getenv("TTS_VOICE", "onyx")
|
||
|
||
SYSTEM_PROMPT = "You are Adolf, a helpful voice assistant. Keep replies concise — 1-3 sentences. No markdown."
|
||
|
||
app = FastAPI(title="Pipecat Voice Bot")
|
||
app.mount("/static", StaticFiles(directory="static"), name="static")
|
||
|
||
|
||
# ── LiveKit helpers ───────────────────────────────────────────────────────────
|
||
def _lk_token(room: str, identity: str, is_bot: bool = False) -> str:
|
||
grants = lkapi.VideoGrants(
|
||
room_join=True,
|
||
room=room,
|
||
can_publish=True,
|
||
can_subscribe=True,
|
||
can_publish_data=True,
|
||
)
|
||
token = (
|
||
lkapi.AccessToken(LIVEKIT_API_KEY, LIVEKIT_SECRET)
|
||
.with_identity(identity)
|
||
.with_name("Adolf Bot" if is_bot else identity)
|
||
.with_grants(grants)
|
||
)
|
||
return token.to_jwt()
|
||
|
||
|
||
async def _create_room(room_name: str) -> None:
|
||
lk = lkapi.LiveKitAPI(LIVEKIT_URL, LIVEKIT_API_KEY, LIVEKIT_SECRET)
|
||
try:
|
||
await lk.room.create_room(
|
||
lkapi.CreateRoomRequest(name=room_name, empty_timeout=300, max_participants=5)
|
||
)
|
||
finally:
|
||
await lk.aclose()
|
||
|
||
|
||
# ── Pipecat pipeline ──────────────────────────────────────────────────────────
|
||
async def _run_bot(room_name: str) -> None:
|
||
bot_token = _lk_token(room_name, "pipecat-bot", is_bot=True)
|
||
|
||
transport = LiveKitTransport(
|
||
url=LIVEKIT_URL,
|
||
token=bot_token,
|
||
room_name=room_name,
|
||
params=LiveKitParams(
|
||
audio_in_enabled=True,
|
||
audio_out_enabled=True,
|
||
vad_enabled=True,
|
||
vad_analyzer=SileroVADAnalyzer(params=VADParams(
|
||
stop_secs=0.8, # wait 0.8s of silence before end-of-speech
|
||
start_secs=0.2, # start speech detection after 0.2s
|
||
confidence=0.85, # high confidence to avoid triggering on ambient noise
|
||
)),
|
||
),
|
||
)
|
||
|
||
stt = OpenAISTTService(
|
||
api_key="dummy",
|
||
base_url=STT_URL,
|
||
model=STT_MODEL,
|
||
language="ru",
|
||
)
|
||
|
||
llm = OpenAILLMService(
|
||
api_key="dummy",
|
||
base_url=ADOLF_URL,
|
||
model="adolf-light",
|
||
)
|
||
|
||
tts = OpenAITTSService(
|
||
api_key="dummy",
|
||
base_url=TTS_URL,
|
||
model="silero",
|
||
voice=TTS_VOICE,
|
||
)
|
||
|
||
messages = [{"role": "system", "content": SYSTEM_PROMPT}]
|
||
context = OpenAILLMContext(messages)
|
||
context_aggregator = llm.create_context_aggregator(context)
|
||
|
||
normalizer = TTSTextNormalizer()
|
||
|
||
pipeline = Pipeline([
|
||
transport.input(),
|
||
stt,
|
||
context_aggregator.user(),
|
||
llm,
|
||
normalizer,
|
||
tts,
|
||
transport.output(),
|
||
context_aggregator.assistant(),
|
||
])
|
||
|
||
task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=False))
|
||
|
||
@transport.event_handler("on_participant_disconnected")
|
||
async def on_disconnect(transport, participant):
|
||
identity = participant if isinstance(participant, str) else getattr(participant, "identity", str(participant))
|
||
logger.info(f"Participant {identity} left — stopping bot")
|
||
await task.cancel()
|
||
|
||
runner = PipelineRunner()
|
||
logger.info(f"Bot starting in room={room_name}")
|
||
await runner.run(task)
|
||
logger.info(f"Bot done in room={room_name}")
|
||
|
||
|
||
# ── API ───────────────────────────────────────────────────────────────────────
|
||
class ConnectResponse(BaseModel):
|
||
room: str
|
||
token: str
|
||
url: str
|
||
|
||
|
||
@app.post("/connect", response_model=ConnectResponse)
|
||
async def connect():
|
||
room_name = f"voice-{uuid.uuid4().hex[:6]}"
|
||
await _create_room(room_name)
|
||
user_token = _lk_token(room_name, "user")
|
||
asyncio.create_task(_run_bot(room_name))
|
||
return ConnectResponse(room=room_name, token=user_token, url=LIVEKIT_PUBLIC_URL)
|
||
|
||
|
||
@app.get("/health")
|
||
async def health():
|
||
return {"status": "ok"}
|
||
|
||
|
||
@app.get("/", response_class=HTMLResponse)
|
||
async def index():
|
||
with open("static/index.html") as f:
|
||
return f.read()
|