from __future__ import annotations import asyncio import base64 import json import logging import os import random import re import time from contextlib import asynccontextmanager from dataclasses import dataclass from typing import Any from urllib.parse import urlencode import websockets from dotenv import load_dotenv from fastapi import FastAPI, WebSocket, WebSocketDisconnect from fastapi.responses import FileResponse from fastapi.staticfiles import StaticFiles import httpx from openai import AsyncOpenAI from app.knowledge import KnowledgeBase, format_knowledge_context, resolve_knowledge_dir load_dotenv() BASE_DIR = os.path.dirname(os.path.abspath(__file__)) PROJECT_DIR = os.path.dirname(BASE_DIR) STATIC_DIR = os.path.join(BASE_DIR, "static") logger = logging.getLogger("simple_call_center") app = FastAPI(title="Simple Call Center") app.mount("/static", StaticFiles(directory=STATIC_DIR), name="static") DEDUPLICATE_TRANSCRIPT_WINDOW_SECS = 4.0 STREAMING_TTS_UNSUPPORTED_MODELS = {"eleven_v3"} UNSUPPORTED_TTS_LANGUAGE_OVERRIDES = { "eleven_flash_v2_5": {"kk"}, "eleven_turbo_v2_5": {"kk"}, } NON_WORD_RE = re.compile(r"[^\w\s]+", re.UNICODE) WHITESPACE_RE = re.compile(r"\s+") ADDRESS_WITH_SLASH_RE = re.compile(r"\b(\d{1,3})\s*/\s*(\d{1,3})\b") NUMBER_WITH_OPTIONAL_LETTER_RE = re.compile(r"\b(\d{1,3})([A-Za-zА-Яа-яЁё])?\b") REPEATED_GREETING_RE = re.compile( r"^\s*(?:здравствуйте|добрый день|доброе утро|добрый вечер)[!.,\s]+", re.IGNORECASE, ) EMAIL_RE = re.compile(r"\S+@\S+") STRUCTURED_LABEL_RE = re.compile(r"\b(Адрес|Телефон|Электронная почта|Почта)\s*:\s*") ACK_WITH_NAME_RE = re.compile( r"^(Хорошо|Поняла|Спасибо),\s+[A-ZА-ЯЁӘҒҚҢӨҰҮҺІ][A-Za-zА-Яа-яЁёӘәҒғҚқҢңӨөҰұҮүҺһІі-]+[,.]\s*" ) ACK_WITH_CITY_RE = re.compile( r"^(Хорошо|Поняла|Спасибо)\.\s*(?:город\s+)?(?:Астана|Алматы|Шымкент),\s*", re.IGNORECASE, ) INITIAL_DETAILS_RE = re.compile( r"\b(?:меня\s+зовут|зовут\s+меня|мое\s+имя|моё\s+имя)\b", re.IGNORECASE, ) CITY_RE = re.compile(r"\b(?:город|г\.)\s*[A-Za-zА-Яа-яЁёӘәҒғҚқҢңӨөҰұҮүҺһІі-]+", re.IGNORECASE) ADDRESS_RE = re.compile(r"\b(?:адрес|улица|ул\.|проспект|пр\.|микрорайон|мкр|дом)\b", re.IGNORECASE) SERVICE_TOPIC_RE = re.compile( r"\b(?:счетчик|счётчик|оплат|начисл|авари|запах|утеч|подключ|отключ|" r"квитанц|долг|лицев|показани|не\s+работает|сломал|сломалась|шумит)\b", re.IGNORECASE, ) BRANCH_QUERY_RE = re.compile( r"(?:какой|какая|название|что\s+за).{0,40}филиал|филиал.{0,40}(?:работает|обслуживает|относится)", re.IGNORECASE, ) CONTACT_REQUEST_RE = re.compile(r"\b(?:адрес|телефон|номер|почт|контакт|куда\s+обращ)\b", re.IGNORECASE) ADDRESS_EXTRACT_RE = re.compile(r"\bадрес\s*[:,-]?\s*(.+)$", re.IGNORECASE) KAZAKH_ADDRESS_EXTRACT_RE = re.compile( r"\b(?:мекенжайым|мекенжай|мекен-жайым|мекен-жай)\s*[:,-]?\s*(.+)$", re.IGNORECASE, ) KAZAKH_SPECIFIC_RE = re.compile(r"[ӘәҒғҚқҢңӨөҰұҮүҺһІі]") KAZAKH_WORD_RE = re.compile( r"\b(?:қазақша|казахша|сәлем|салем|аты(?:м)?|есім(?:ім)?|қала(?:м|сы)?|" r"мекенжай|мекен-жай|көмек|төлем|есептегіш|көрсеткіш)\b", re.IGNORECASE, ) KAZAKH_LANGUAGE_SELECTION_RE = re.compile( r"\b(?:қазақша|казахша|қазақ|казахский|казахском|kk)\b", re.IGNORECASE, ) RUSSIAN_LANGUAGE_SELECTION_RE = re.compile( r"\b(?:орысша|русский|русском|по-русски|русски|ru)\b", re.IGNORECASE, ) LANGUAGE_SELECTION_WORDS_RE = re.compile( r"\b(?:қазақша|казахша|қазақ|казахский|казахском|орысша|русский|русском|" r"по-русски|русски|kk|ru)\b", re.IGNORECASE, ) RUSSIAN_WORD_RE = re.compile( r"\b(?:здравствуйте|город|адрес|улица|подскажите|пожалуйста|помочь|" r"оплата|счетчик|счётчик|начисления|авария)\b", re.IGNORECASE, ) KAZAKH_INITIAL_DETAILS_RE = re.compile( r"\b(?:менің\s+атым|атым|есімім|мені\s+.*?деп\s+атайды)\b", re.IGNORECASE, ) KAZAKH_CITY_RE = re.compile(r"\b(?:қала(?:м|сы)?|қ\.)\b", re.IGNORECASE) KAZAKH_ADDRESS_RE = re.compile( r"\b(?:мекенжай(?:ым)?|мекен-жай(?:ым)?|көше|көшесі|үй|адрес)\b", re.IGNORECASE, ) FILLER_PREFIX_RE = re.compile(r"^\s*(?:м-м|ммм|эм|эмм|эм-м)[,.\s]+", re.IGNORECASE) FILLER_BLOCK_RE = re.compile( r"\b(?:авари|утеч|запах|опасн|срочно|один\s+ноль\s+четыре|телефон|номер|почт|" r"до\s+свидания|извин)\b", re.IGNORECASE, ) CYRILLIC_TITLE_START_RE = re.compile(r"^([А-ЯЁӘҒҚҢӨҰҮҺІ])([а-яёәғқңөұүһі])") ONES = { 0: "ноль", 1: "один", 2: "два", 3: "три", 4: "четыре", 5: "пять", 6: "шесть", 7: "семь", 8: "восемь", 9: "девять", } TEENS = { 10: "десять", 11: "одиннадцать", 12: "двенадцать", 13: "тринадцать", 14: "четырнадцать", 15: "пятнадцать", 16: "шестнадцать", 17: "семнадцать", 18: "восемнадцать", 19: "девятнадцать", } TENS = { 20: "двадцать", 30: "тридцать", 40: "сорок", 50: "пятьдесят", 60: "шестьдесят", 70: "семьдесят", 80: "восемьдесят", 90: "девяносто", } HUNDREDS = { 100: "сто", 200: "двести", 300: "триста", 400: "четыреста", 500: "пятьсот", 600: "шестьсот", 700: "семьсот", 800: "восемьсот", 900: "девятьсот", } KAZAKH_ONES = { 0: "нөл", 1: "бір", 2: "екі", 3: "үш", 4: "төрт", 5: "бес", 6: "алты", 7: "жеті", 8: "сегіз", 9: "тоғыз", } KAZAKH_TENS = { 10: "он", 20: "жиырма", 30: "отыз", 40: "қырық", 50: "елу", 60: "алпыс", 70: "жетпіс", 80: "сексен", 90: "тоқсан", } KAZAKH_HUNDREDS = { 100: "жүз", 200: "екі жүз", 300: "үш жүз", 400: "төрт жүз", 500: "бес жүз", 600: "алты жүз", 700: "жеті жүз", 800: "сегіз жүз", 900: "тоғыз жүз", } DEFAULT_INITIAL_GREETING = ( "Сәлеметсіз бе! Операторға дұрыс бағыттау үшін қызмет көрсету тілін таңдаңыз: " "қазақша немесе орысша. Здравствуйте! Чтобы правильно перевести вас на оператора, " "выберите язык обслуживания: казахский или русский" ) KB_INSTRUCTIONS = ( "После того как имя, город и адрес клиента собраны, используй базу знаний QazaqGasAimaq " "для ответов по услугам, регламентам, аварийным ситуациям, оплате, начислениям, " "подключению, отключению, приборам учета и контактам филиалов. Если в текущем " "запросе есть служебный контекст базы знаний, опирайся на него и отвечай кратко. " "Главная задача — помочь клиенту самостоятельно в рамках диалога. Не предлагай " "подключить, перевести или передать клиента специалисту, если клиент сам этого " "явно не попросил. Контакты филиалов выдавай только когда клиент прямо спрашивает " "адрес, телефон, филиал или куда обращаться. Не перечисляй адрес, почту и телефон " "одновременно, если это не нужно. Для голосового ответа сначала дай короткий " "человеческий вывод, например: «По Астане это Астанинский филиал». Электронную " "почту называй только если клиент сам попросил почту. Если релевантного контекста " "нет или ответа в базе не хватает, " "не выдумывай: задай один короткий уточняющий вопрос или предложи ближайший " "безопасный следующий шаг по доступной информации." ) VOICE_READY_INSTRUCTIONS = ( "Все ответы должны быть готовы для озвучки. Не используй цифры в клиентском ответе. " "Числа, сроки, суммы, адресные номера и телефоны пиши словами так, как их нужно " "произнести. Телефонные номера произноси словами по цифрам или удобными короткими " "группами, без символов плюс, скобок и дефисов. Аварийный номер сто четыре всегда " "пиши как «один ноль четыре». Адресные дроби и номера со слешем читай как «дробь»: " "«дом два дробь два», а не как математическое деление. Буквы в адресах произноси " "как буквы: «сорок шесть Б»." ) HUMAN_RESPONSE_GUARDRAILS = ( "Критически важные правила живой речи. По умолчанию не используй имя клиента в " "ответах: имя можно произнести только если нужно переспросить его или клиент сам " "просит обращаться по имени. В первой реплике после получения имени, города и " "адреса не обращайся по имени, не повторяй город и не делай пустую строку. Не " "перечисляй варианты обращения как меню через запятую: «счетчик, оплата, " "начисления, подключение, аварийная ситуация». Сначала задай один открытый вопрос. " "Правильная форма после сбора данных: «Хорошо. С чем помочь по адресу Туран сорок " "шесть Б?» или «Поняла, Туран сорок шесть Б. Что случилось?». Неправильно: " "«Хорошо, Магжан» и неправильно: «Например: счетчик, оплата, начисления или " "аварийная ситуация»." ) DEFAULT_OPERATOR_STYLE_INSTRUCTIONS = ( "Говори как живой оператор в коротком телефонном разговоре. Реплики обычно одна-две " "короткие фразы. Не отвечай справочной карточкой: без заголовков «Адрес», «Телефон», " "«Электронная почта», без списков и без перечисления всех контактов подряд. Не начинай " "каждую реплику со слов «спасибо», «поняла», «отлично», «зафиксировала» или имени " "клиента. Слово «отлично» не используй для данных клиента или проблемной ситуации. " "Не говори «адрес зафиксировала», «данные зафиксировала», «для CRM», «передам " "специалисту», «подключу специалиста», если клиент сам не просил. Если имя могло " "быть распознано неверно, не обращайся по имени; используй нейтральное обращение. " "Даже когда имя известно, используй его редко и не используй имя в ответе сразу " "после сбора первичных данных. После получения имени, города и адреса не пересказывай " "их полностью и не повторяй город рядом с адресом. Хороший ответ: «Хорошо. С чем " "помочь по адресу Туран сорок шесть Б?» или «Поняла, Туран сорок шесть Б. Что " "случилось?». Плохой ответ: «Отлично, Магжан. Адрес зафиксировала: город Астана, " "улица Туран, дом сорок шесть Б». Благодари только когда это звучит уместно. " "В обычном диалоге используй живые фразы: «сейчас посмотрю», «давайте проверим», " "«подскажу», «уточню один момент». Избегай канцелярита, роботичных подтверждений " "и фраз вроде «вам удобнее получить информацию»." ) @dataclass(frozen=True) class Settings: openai_api_key: str | None = os.getenv("OPENAI_API_KEY") elevenlabs_api_key: str | None = os.getenv("ELEVENLABS_API_KEY") openai_model: str = os.getenv("OPENAI_MODEL", "gpt-5.4-nano") openai_reasoning_effort: str = os.getenv("OPENAI_REASONING_EFFORT", "none") tts_provider: str = os.getenv("AI_VOICE_TTS_PROVIDER", "elevenlabs") streaming_tts_enabled: bool = os.getenv("AI_VOICE_V2_STREAMING_TTS", "1") == "1" elevenlabs_voice_id: str = ( os.getenv("AI_VOICE_TTS_ELEVENLABS_VOICE_ID") or os.getenv("ELEVENLABS_VOICE_ID", "4O1sYUnmtThcBoSBrri7") ) elevenlabs_stt_model: str = os.getenv("ELEVENLABS_STT_MODEL", "scribe_v2_realtime") elevenlabs_requested_tts_model: str = ( os.getenv("AI_VOICE_TTS_ELEVENLABS_MODEL_ID") or os.getenv("ELEVENLABS_TTS_MODEL", "eleven_v3") ) elevenlabs_tts_output_format: str = ( os.getenv("AI_VOICE_TTS_ELEVENLABS_OUTPUT_FORMAT") or os.getenv("ELEVENLABS_TTS_OUTPUT_FORMAT", "pcm_16000") ) elevenlabs_streaming_fallback_model: str = os.getenv( "AI_VOICE_TTS_ELEVENLABS_STREAMING_FALLBACK_MODEL", "eleven_flash_v2_5", ) elevenlabs_tts_speed: float = float(os.getenv("ELEVENLABS_TTS_SPEED", "1.0")) elevenlabs_tts_stability: float = float(os.getenv("ELEVENLABS_TTS_STABILITY", "0.4")) elevenlabs_tts_similarity_boost: float = float( os.getenv("ELEVENLABS_TTS_SIMILARITY_BOOST", "0.8") ) elevenlabs_stt_language: str = os.getenv("ELEVENLABS_STT_LANGUAGE", "") elevenlabs_tts_language: str = os.getenv("ELEVENLABS_TTS_LANGUAGE", "") initial_greeting_tts_language: str = os.getenv("INITIAL_GREETING_TTS_LANGUAGE", "") default_conversation_language: str = os.getenv("DEFAULT_CONVERSATION_LANGUAGE", "ru") humanizer_enabled: bool = os.getenv("HUMANIZER_ENABLED", "1") == "1" humanizer_filler_rate: float = float(os.getenv("HUMANIZER_FILLER_RATE", "0.15")) humanizer_fillers: tuple[str, ...] = tuple( filler.strip() for filler in os.getenv("HUMANIZER_FILLERS", "м-м,эм").split(",") if filler.strip() ) local_kb_enabled: bool = os.getenv("LOCAL_KB_ENABLED", "1") == "1" local_kb_dir: str = resolve_knowledge_dir(PROJECT_DIR, os.getenv("LOCAL_KB_DIR", "knowledge")) local_kb_max_results: int = int(os.getenv("LOCAL_KB_MAX_RESULTS", "4")) local_kb_min_score: float = float(os.getenv("LOCAL_KB_MIN_SCORE", "2.0")) initial_greeting: str = os.getenv("INITIAL_GREETING", DEFAULT_INITIAL_GREETING) operator_style_instructions: str = os.getenv( "OPERATOR_STYLE_INSTRUCTIONS", DEFAULT_OPERATOR_STYLE_INSTRUCTIONS, ) assistant_instructions: str = os.getenv( "ASSISTANT_INSTRUCTIONS", "Ты Айнур, оператор контакт-центра QazaqGasAimaq для входящих звонков. " "Твоя задача - коротко собрать первичные данные клиента для обработки обращения. " "Обязательные первичные данные: имя клиента, город и адрес. В начале разговора " "обязательно попроси недостающие данные: если имени еще нет, спроси имя; если " "города или адреса еще нет, спроси «подскажите пожалуйста ваш город и адрес». " "Если клиент озвучил только часть обязательных данных, вежливо попроси только " "то, чего не хватает. Пока имя, город и адрес не получены, не помогай по сути " "обращения и не давай консультацию; коротко объясни, что сначала нужно " "уточнить данные для обработки обращения. После ответа клиента не " "спрашивай уже полученные данные повторно; для справочных вопросов считай " "адрес достаточным, если клиент назвал город и понятный адрес, улицу или дом. " "Не требуй квартиру или полный адрес, если это не нужно для конкретной заявки. " "Используй полученные данные как внутренний контекст, но не проговаривай клиенту " "слова CRM, аналитика или «зафиксировала». Отвечай естественно, кратко и по делу. " "Сначала старайся сама обработать вопрос и помочь клиенту по доступной " "информации; не предлагай подключить, перевести или передать клиента " "специалисту, если клиент сам этого явно не попросил. " "Когда говоришь о себе, используй женский род: могла, смогла, сделала, готова, " "проверила, нашла. Не завершай диалог самостоятельно и говори «до свидания» " "только если клиент явно попрощался или попросил завершить разговор. Все ответы " "должны быть готовы для озвучки: не используй цифры в клиентском ответе. Числа, " "сроки, суммы, адресные номера и телефоны пиши словами так, как их нужно " "произнести. Телефонные номера произноси словами по цифрам или удобными короткими " "группами, без символов плюс, скобок и дефисов. Не используй Markdown, URL или таблицы", ) @property def elevenlabs_tts_model(self) -> str: if ( self.streaming_tts_enabled and self.elevenlabs_requested_tts_model in STREAMING_TTS_UNSUPPORTED_MODELS ): return self.elevenlabs_streaming_fallback_model return self.elevenlabs_requested_tts_model @property def tts_warning(self) -> str | None: if self.elevenlabs_tts_model == self.elevenlabs_requested_tts_model: return None return ( f"{self.elevenlabs_requested_tts_model} is not supported by ElevenLabs " f"WebSocket streaming; using {self.elevenlabs_tts_model} instead." ) settings = Settings() openai_client = AsyncOpenAI(api_key=settings.openai_api_key) if settings.openai_api_key else None knowledge_base = ( KnowledgeBase.from_directory(settings.local_kb_dir) if settings.local_kb_enabled else KnowledgeBase.empty() ) if settings.tts_warning: logger.warning("TTS fallback active: %s", settings.tts_warning) logger.info( "Local knowledge base: enabled=%s documents=%s dir=%s", settings.local_kb_enabled, knowledge_base.size, settings.local_kb_dir, ) @app.get("/") async def index() -> FileResponse: return FileResponse(os.path.join(STATIC_DIR, "index.html")) @app.get("/health") async def health() -> dict[str, Any]: return { "ok": True, "openai_key": bool(settings.openai_api_key), "elevenlabs_key": bool(settings.elevenlabs_api_key), "openai_model": settings.openai_model, "tts_provider": settings.tts_provider, "streaming_tts_enabled": settings.streaming_tts_enabled, "stt_model": settings.elevenlabs_stt_model, "requested_tts_model": settings.elevenlabs_requested_tts_model, "tts_model": settings.elevenlabs_tts_model, "tts_output_format": settings.elevenlabs_tts_output_format, "voice_id": settings.elevenlabs_voice_id, "tts_language": settings.elevenlabs_tts_language or None, "initial_greeting_tts_language": settings.initial_greeting_tts_language or None, "default_conversation_language": settings.default_conversation_language, "tts_speed": settings.elevenlabs_tts_speed, "tts_warning": settings.tts_warning, "humanizer_enabled": settings.humanizer_enabled, "humanizer_filler_rate": settings.humanizer_filler_rate, "humanizer_fillers": list(settings.humanizer_fillers), "local_kb_enabled": settings.local_kb_enabled, "local_kb_documents": knowledge_base.size, "local_kb_dir": settings.local_kb_dir, "initial_greeting": settings.initial_greeting, "operator_style_preview": settings.operator_style_instructions[:80], "assistant_prompt_preview": settings.assistant_instructions[:80], } @asynccontextmanager async def connect_ws(uri: str, headers: dict[str, str] | None = None): try: websocket_cm = websockets.connect(uri, additional_headers=headers, max_size=None) except TypeError: websocket_cm = websockets.connect(uri, extra_headers=headers, max_size=None) async with websocket_cm as websocket: yield websocket class VoiceSession: def __init__(self, client_ws: WebSocket): self.client_ws = client_ws self.history: list[dict[str, str]] = [] self.user_text_queue: asyncio.Queue[str | None] = asyncio.Queue() self.send_lock = asyncio.Lock() self.active_user_text_keys: set[str] = set() self.recent_user_text_keys: dict[str, float] = {} self.initial_greeting_sent = False self.responding = False self.closed = False self.conversation_language: str | None = None self.awaiting_language_selection = False async def send_json(self, payload: dict[str, Any]) -> None: if self.closed: return async with self.send_lock: await self.client_ws.send_json(payload) async def send_tts_failed(self, message: str) -> None: logger.warning("TTS failed: %s", message) await self.send_json({"type": "tts_failed", "message": message}) async def run(self) -> None: await self.client_ws.accept() missing = [] if not settings.elevenlabs_api_key: missing.append("ELEVENLABS_API_KEY") if not settings.openai_api_key: missing.append("OPENAI_API_KEY") if settings.tts_provider != "elevenlabs": missing.append("AI_VOICE_TTS_PROVIDER=elevenlabs") if missing: await self.send_json( { "type": "error", "message": f"Missing environment variables: {', '.join(missing)}", } ) await self.client_ws.close() return await self.send_json( { "type": "ready", "model": settings.openai_model, "stt_model": settings.elevenlabs_stt_model, "tts_model": settings.elevenlabs_tts_model, "requested_tts_model": settings.elevenlabs_requested_tts_model, "streaming_tts_enabled": settings.streaming_tts_enabled, "tts_warning": settings.tts_warning, } ) stt_uri = self.build_stt_uri() headers = {"xi-api-key": settings.elevenlabs_api_key or ""} async with connect_ws(stt_uri, headers=headers) as stt_ws: tasks = [ asyncio.create_task(self.forward_audio_to_stt(stt_ws)), asyncio.create_task(self.forward_stt_to_client(stt_ws)), asyncio.create_task(self.process_user_texts()), ] done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED) for task in pending: task.cancel() for task in done: task.result() def build_stt_uri(self) -> str: params = { "model_id": settings.elevenlabs_stt_model, "audio_format": "pcm_16000", "commit_strategy": "vad", "vad_silence_threshold_secs": "0.8", "include_language_detection": "true", } if settings.elevenlabs_stt_language: params["language_code"] = settings.elevenlabs_stt_language return f"wss://api.elevenlabs.io/v1/speech-to-text/realtime?{urlencode(params)}" async def forward_audio_to_stt(self, stt_ws: Any) -> None: try: while True: message = await self.client_ws.receive() if message.get("type") == "websocket.disconnect": raise WebSocketDisconnect() audio_bytes = message.get("bytes") if audio_bytes: if self.responding: continue await stt_ws.send( json.dumps( { "message_type": "input_audio_chunk", "audio_base_64": base64.b64encode(audio_bytes).decode("ascii"), "sample_rate": 16000, } ) ) continue text_message = message.get("text") if text_message: await self.handle_client_command(text_message) except (WebSocketDisconnect, RuntimeError): self.closed = True await self.user_text_queue.put(None) async def handle_client_command(self, text_message: str) -> None: try: command = json.loads(text_message) except json.JSONDecodeError: return if command.get("type") == "reset": self.history.clear() self.active_user_text_keys.clear() self.recent_user_text_keys.clear() self.initial_greeting_sent = False self.conversation_language = None self.awaiting_language_selection = False await self.send_json({"type": "reset_done"}) elif command.get("type") == "start_greeting": await self.send_initial_greeting() async def send_initial_greeting(self) -> None: if self.initial_greeting_sent: return self.initial_greeting_sent = True self.awaiting_language_selection = True await self.speak_assistant_text( settings.initial_greeting, add_to_history=True, tts_language=settings.initial_greeting_tts_language, ) async def speak_assistant_text( self, text: str, add_to_history: bool = False, tts_language: str | None = None, ) -> None: self.responding = True await self.send_json({"type": "assistant_started"}) try: await self.send_json({"type": "assistant_delta", "text": text}) await self.send_json({"type": "assistant_text_done", "text": text}) if settings.streaming_tts_enabled: tts_queue: asyncio.Queue[str | None] = asyncio.Queue() tts_task = asyncio.create_task(self.stream_tts(tts_queue, tts_language)) await tts_queue.put(text) await tts_queue.put(None) await tts_task else: await self.synthesize_tts_http(text, tts_language) if add_to_history: self.history.append({"role": "assistant", "content": text}) self.history = self.history[-16:] finally: await self.send_json({"type": "assistant_done"}) self.responding = False async def forward_stt_to_client(self, stt_ws: Any) -> None: async for raw_message in stt_ws: data = json.loads(raw_message) message_type = data.get("message_type") if message_type == "session_started": await self.send_json({"type": "stt_ready"}) elif message_type == "partial_transcript": await self.send_json({"type": "stt_partial", "text": data.get("text", "")}) elif message_type in {"committed_transcript", "committed_transcript_with_timestamps"}: text = (data.get("text") or "").strip() if text and self.accept_committed_transcript(text): language = self.update_conversation_language(text, data.get("language_code")) await self.send_json( { "type": "user_final", "text": text, "language": language, } ) await self.user_text_queue.put(text) elif message_type and "error" in message_type: await self.send_json( { "type": "error", "message": data.get("message") or data.get("error") or raw_message, } ) def update_conversation_language( self, text: str, provider_language: str | None = None, ) -> str: detected = detect_conversation_language(text, provider_language) if detected: self.conversation_language = detected if self.conversation_language: return self.conversation_language return normalize_language_code(settings.default_conversation_language) or "ru" def handle_language_selection(self, user_text: str) -> str | None: if self.conversation_language and not self.awaiting_language_selection: return None selected_language = detect_language_selection(user_text) if selected_language is None: selected_language = detect_conversation_language(user_text) if selected_language in {"kk", "ru"}: self.conversation_language = selected_language self.awaiting_language_selection = False if has_initial_customer_details(user_text): return None return initial_details_prompt_for_language(selected_language) if self.awaiting_language_selection: return language_selection_reprompt() return None def tts_language_for_session(self) -> str: return self.conversation_language or normalize_language_code(settings.elevenlabs_tts_language) or "ru" def language_instruction_for_session(self) -> str: language = self.conversation_language or normalize_language_code(settings.default_conversation_language) if language == "kk": return ( "Клиент выбрал казахский язык или говорит по-казахски. Отвечай только на казахском языке. " "Сохраняй роль Айнур, оператора QazaqGasAimaq. Не смешивай русский и казахский, кроме названия компании. " "Собирай имя, город и мекенжай, если чего-то не хватает. Отвечай қысқа, табиғи және нақты." " Числа, телефондар және мекенжай нөмірлерін қазақша сөзбен жаз." ) return ( "Клиент говорит по-русски. Отвечай только на русском языке, если клиент явно не перешел на казахский." ) def accept_committed_transcript(self, text: str) -> bool: key = normalize_transcript(text) if not key: return False now = time.monotonic() self.expire_recent_user_text_keys(now) if key in self.active_user_text_keys or key in self.recent_user_text_keys: return False self.active_user_text_keys.add(key) return True def mark_committed_transcript_done(self, text: str) -> None: key = normalize_transcript(text) if not key: return now = time.monotonic() self.active_user_text_keys.discard(key) self.recent_user_text_keys[key] = now self.expire_recent_user_text_keys(now) def expire_recent_user_text_keys(self, now: float) -> None: expired_keys = [ key for key, seen_at in self.recent_user_text_keys.items() if now - seen_at > DEDUPLICATE_TRANSCRIPT_WINDOW_SECS ] for key in expired_keys: del self.recent_user_text_keys[key] async def process_user_texts(self) -> None: while True: text = await self.user_text_queue.get() if text is None: return try: await self.respond_to_user(text) finally: self.mark_committed_transcript_done(text) async def respond_to_user(self, user_text: str) -> None: if openai_client is None: await self.send_json({"type": "error", "message": "OpenAI client is not configured."}) return self.responding = True await self.send_json({"type": "assistant_started"}) language_selection_answer = self.handle_language_selection(user_text) if not language_selection_answer: self.update_conversation_language(user_text) tts_language = self.tts_language_for_session() tts_queue: asyncio.Queue[str | None] | None = None tts_task: asyncio.Task[None] | None = None if settings.streaming_tts_enabled: tts_queue = asyncio.Queue() tts_task = asyncio.create_task(self.stream_tts(tts_queue, tts_language)) answer_parts: list[str] = [] try: if language_selection_answer: await self.emit_assistant_answer( user_text, language_selection_answer, tts_queue, tts_language, ) return quick_answer = build_initial_details_ack(user_text, self.conversation_language) if quick_answer: await self.emit_assistant_answer(user_text, quick_answer, tts_queue, tts_language) return retrieval_query = " ".join( message["content"] for message in self.history[-6:] if message.get("role") == "user" ) retrieval_query = f"{retrieval_query} {user_text}".strip() kb_hits = ( knowledge_base.search( retrieval_query, limit=settings.local_kb_max_results, min_score=settings.local_kb_min_score, ) if settings.local_kb_enabled else [] ) kb_context = format_knowledge_context(kb_hits) kb_context = prepare_model_context_text(kb_context) if kb_hits: logger.info( "KB hits for %r: %s", user_text, ", ".join( f"{hit.record.get('external_id')}:{hit.score:.1f}" for hit in kb_hits ), ) effective_user_text = user_text if kb_context: effective_user_text = ( f"{kb_context}\n\n" f"Текущая реплика клиента, на которую нужно ответить: {user_text}" ) input_messages = [*self.history, {"role": "user", "content": effective_user_text}] request: dict[str, Any] = { "model": settings.openai_model, "instructions": ( f"{settings.assistant_instructions}\n\n" f"{self.language_instruction_for_session()}\n\n" f"{KB_INSTRUCTIONS}\n\n" f"{settings.operator_style_instructions}\n\n" f"{VOICE_READY_INSTRUCTIONS}\n\n" f"{HUMAN_RESPONSE_GUARDRAILS}" ), "input": input_messages, "stream": True, "store": False, "max_output_tokens": 500, } if settings.openai_reasoning_effort: request["reasoning"] = {"effort": settings.openai_reasoning_effort} stream = await openai_client.responses.create(**request) async for event in stream: event_type = getattr(event, "type", "") if event_type == "response.output_text.delta": delta = getattr(event, "delta", "") if delta: answer_parts.append(delta) elif event_type == "error": await self.send_json( {"type": "error", "message": getattr(event, "message", "OpenAI stream error")} ) answer = polish_assistant_answer("".join(answer_parts), user_text, self.conversation_language) answer = humanize_assistant_answer(answer, user_text, self.conversation_language) await self.emit_assistant_answer(user_text, answer, tts_queue, tts_language) except Exception as exc: await self.send_json({"type": "error", "message": f"OpenAI error: {exc}"}) finally: if tts_queue is not None: await tts_queue.put(None) if tts_task is not None: await tts_task await self.send_json({"type": "assistant_done"}) self.responding = False async def emit_assistant_answer( self, user_text: str, answer: str, tts_queue: asyncio.Queue[str | None] | None, tts_language: str | None = None, ) -> None: if answer: await self.send_json({"type": "assistant_delta", "text": answer}) if tts_queue is not None: await tts_queue.put(answer) self.history.extend( [ {"role": "user", "content": user_text}, {"role": "assistant", "content": answer}, ] ) self.history = self.history[-16:] await self.send_json({"type": "assistant_text_done", "text": answer}) if answer and not settings.streaming_tts_enabled: await self.synthesize_tts_http(answer, tts_language) async def synthesize_tts_http(self, text: str, language_code: str | None = None) -> None: voice_id = settings.elevenlabs_voice_id tts_language = resolve_tts_language(language_code) logger.info( "Calling ElevenLabs HTTP TTS: model=%s voice_id=%s language=%s", settings.elevenlabs_tts_model, voice_id, tts_language or "auto", ) params = {"output_format": settings.elevenlabs_tts_output_format} payload: dict[str, Any] = { "text": prepare_tts_text(text, tts_language), "model_id": settings.elevenlabs_tts_model, "voice_settings": { "stability": settings.elevenlabs_tts_stability, "similarity_boost": settings.elevenlabs_tts_similarity_boost, "use_speaker_boost": False, "speed": settings.elevenlabs_tts_speed, }, } if tts_language: payload["language_code"] = tts_language try: async with httpx.AsyncClient(timeout=60.0) as client: response = await client.post( f"https://api.elevenlabs.io/v1/text-to-speech/{voice_id}", params=params, headers={ "xi-api-key": settings.elevenlabs_api_key or "", "Content-Type": "application/json", }, json=payload, ) response.raise_for_status() await self.send_json( { "type": "tts_audio", "audio": base64.b64encode(response.content).decode("ascii"), "format": audio_format_label(settings.elevenlabs_tts_output_format), "sample_rate": audio_sample_rate(settings.elevenlabs_tts_output_format), "mime_type": audio_mime_type(settings.elevenlabs_tts_output_format), } ) except httpx.HTTPStatusError as exc: message = exc.response.text[:500] if exc.response is not None else str(exc) await self.send_tts_failed( f"ElevenLabs HTTP TTS error: {exc.response.status_code} {message}" ) except Exception as exc: await self.send_tts_failed(f"ElevenLabs HTTP TTS error: {exc}") async def stream_tts( self, text_queue: asyncio.Queue[str | None], language_code: str | None = None, ) -> None: voice_id = settings.elevenlabs_voice_id tts_language = resolve_tts_language(language_code) logger.info( "Opening ElevenLabs TTS websocket: requested_model=%s effective_model=%s voice_id=%s language=%s", settings.elevenlabs_requested_tts_model, settings.elevenlabs_tts_model, voice_id, tts_language or "auto", ) params = { "model_id": settings.elevenlabs_tts_model, "output_format": settings.elevenlabs_tts_output_format, "inactivity_timeout": "180", } if tts_language: params["language_code"] = tts_language uri = ( f"wss://api.elevenlabs.io/v1/text-to-speech/{voice_id}/stream-input?" f"{urlencode(params)}" ) headers = {"xi-api-key": settings.elevenlabs_api_key or ""} try: async with connect_ws(uri, headers=headers) as tts_ws: await tts_ws.send( json.dumps( { "text": " ", "voice_settings": { "stability": settings.elevenlabs_tts_stability, "similarity_boost": settings.elevenlabs_tts_similarity_boost, "use_speaker_boost": False, "speed": settings.elevenlabs_tts_speed, }, "generation_config": {"chunk_length_schedule": [50, 80, 120, 160]}, "xi_api_key": settings.elevenlabs_api_key, } ) ) sender_done = asyncio.Event() sender_task = asyncio.create_task( self.send_text_to_tts(tts_ws, text_queue, sender_done, tts_language) ) receiver_task = asyncio.create_task(self.receive_tts_audio(tts_ws, sender_done)) try: await sender_task await asyncio.wait_for(receiver_task, timeout=30.0) except asyncio.TimeoutError: logger.warning("ElevenLabs TTS receiver timed out after text sender completed.") finally: for task in (sender_task, receiver_task): if not task.done(): task.cancel() except Exception as exc: await self.send_tts_failed(f"ElevenLabs TTS error: {exc}") async def send_text_to_tts( self, tts_ws: Any, text_queue: asyncio.Queue[str | None], sender_done: asyncio.Event, language_code: str | None = None, ) -> None: buffer = "" try: while True: chunk = await text_queue.get() if chunk is None: if buffer: await tts_ws.send( json.dumps( {"text": prepare_tts_text(buffer, language_code), "flush": True} ) ) await tts_ws.send(json.dumps({"text": ""})) return buffer += chunk if len(buffer) >= 90 or any(buffer.endswith(mark) for mark in ".!?;:\n"): await tts_ws.send(json.dumps({"text": prepare_tts_text(buffer, language_code)})) buffer = "" finally: sender_done.set() async def receive_tts_audio(self, tts_ws: Any, sender_done: asyncio.Event) -> None: async for raw_message in tts_ws: data = json.loads(raw_message) audio = data.get("audio") if audio: await self.send_json( { "type": "tts_audio", "audio": audio, "format": audio_format_label(settings.elevenlabs_tts_output_format), "sample_rate": audio_sample_rate(settings.elevenlabs_tts_output_format), "mime_type": audio_mime_type(settings.elevenlabs_tts_output_format), } ) if (data.get("isFinal") or data.get("is_final")) and sender_done.is_set(): return def normalize_transcript(text: str) -> str: text = NON_WORD_RE.sub("", text) text = WHITESPACE_RE.sub(" ", text) return text.strip().casefold() def normalize_language_code(language_code: str | None) -> str | None: if not language_code: return None normalized = language_code.strip().casefold().replace("_", "-") if not normalized: return None if normalized in {"kk", "kz", "kaz", "kazakh"} or normalized.startswith("kk-"): return "kk" if normalized in {"ru", "rus", "russian"} or normalized.startswith("ru-"): return "ru" if normalized in {"en", "eng", "english"} or normalized.startswith("en-"): return "en" return normalized.split("-", 1)[0] def resolve_tts_language(language_code: str | None) -> str | None: if language_code is not None: raw = language_code.strip().casefold() if not raw or raw in {"auto", "none", "detect"}: return None return supported_tts_language_or_auto(normalize_language_code(raw)) normalized = normalize_language_code(language_code) if normalized: return supported_tts_language_or_auto(normalized) return supported_tts_language_or_auto(normalize_language_code(settings.elevenlabs_tts_language)) def supported_tts_language_or_auto(language_code: str | None) -> str | None: if not language_code: return None unsupported = UNSUPPORTED_TTS_LANGUAGE_OVERRIDES.get(settings.elevenlabs_tts_model, set()) if language_code in unsupported: return None return language_code def detect_conversation_language(text: str, provider_language: str | None = None) -> str | None: provider = normalize_language_code(provider_language) if KAZAKH_SPECIFIC_RE.search(text) or KAZAKH_WORD_RE.search(text): return "kk" if RUSSIAN_WORD_RE.search(text): return "ru" if provider in {"kk", "ru"}: return provider return provider if provider in {"en"} else None def detect_language_selection(text: str) -> str | None: wants_kazakh = bool(KAZAKH_LANGUAGE_SELECTION_RE.search(text)) wants_russian = bool(RUSSIAN_LANGUAGE_SELECTION_RE.search(text)) if wants_kazakh and not wants_russian: return "kk" if wants_russian and not wants_kazakh: return "ru" return None def has_initial_customer_details(text: str) -> bool: return bool( INITIAL_DETAILS_RE.search(text) or KAZAKH_INITIAL_DETAILS_RE.search(text) or CITY_RE.search(text) or KAZAKH_CITY_RE.search(text) or ADDRESS_RE.search(text) or KAZAKH_ADDRESS_RE.search(text) ) def initial_details_prompt_for_language(language_code: str) -> str: if normalize_language_code(language_code) == "kk": return ( "Жақсы, қазақша жалғастырамыз. Көмектесуім үшін атыңызды, қалаңызды " "және мекенжайыңызды айтыңызшы." ) return ( "Хорошо, продолжим на русском. Чтобы я смогла помочь, подскажите, пожалуйста, " "ваше имя, город и адрес." ) def language_selection_reprompt() -> str: return ( "Қызмет көрсету тілін таңдаңыз: қазақша немесе орысша. " "Выберите язык обслуживания: казахский или русский." ) def build_initial_details_ack(user_text: str, language_code: str | None = None) -> str | None: language = normalize_language_code(language_code) or detect_conversation_language(user_text) if language == "kk": return build_kazakh_initial_details_ack(user_text) if not INITIAL_DETAILS_RE.search(user_text): return None if not CITY_RE.search(user_text) or not ADDRESS_RE.search(user_text): return None if "?" in user_text or SERVICE_TOPIC_RE.search(user_text): return None address = extract_address_fragment(user_text) if address: return f"Хорошо. С чем помочь по адресу {voice_ready_address(address)}?" return "Хорошо. С чем помочь по этому адресу?" def build_kazakh_initial_details_ack(user_text: str) -> str | None: if not KAZAKH_INITIAL_DETAILS_RE.search(user_text): return None if not KAZAKH_CITY_RE.search(user_text) or not KAZAKH_ADDRESS_RE.search(user_text): return None if "?" in user_text or SERVICE_TOPIC_RE.search(user_text): return None address = extract_kazakh_address_fragment(user_text) if address: return f"Жақсы. {voice_ready_address(address, 'kk')} мекенжайы бойынша немен көмектесейін?" return "Жақсы. Осы мекенжай бойынша немен көмектесейін?" def extract_address_fragment(text: str) -> str: match = ADDRESS_EXTRACT_RE.search(text) if not match: return "" address = match.group(1) address = re.split( r"\b(?:у\s+меня|мне\s+нужно|хочу|подскажите|скажите|вопрос|проблема)\b", address, maxsplit=1, flags=re.IGNORECASE, )[0] return address.strip(" .,!?:;") def extract_kazakh_address_fragment(text: str) -> str: match = KAZAKH_ADDRESS_EXTRACT_RE.search(text) if not match: return "" address = match.group(1) address = re.split( r"\b(?:маған|менде|сұрақ|мәселе|көмек|керек|қажет)\b", address, maxsplit=1, flags=re.IGNORECASE, )[0] return address.strip(" .,!?:;") def prepare_model_context_text(text: str) -> str: return text def polish_assistant_answer( answer: str, user_text: str, language_code: str | None = None, ) -> str: answer = WHITESPACE_RE.sub(" ", answer).strip() if not answer: return "" language = normalize_language_code(language_code) or detect_conversation_language(user_text) or "ru" quick_answer = build_initial_details_ack(user_text, language) if quick_answer: return quick_answer if language == "ru": answer = REPEATED_GREETING_RE.sub("", answer).strip() answer = ACK_WITH_NAME_RE.sub(r"\1. ", answer).strip() answer = ACK_WITH_CITY_RE.sub(r"\1. ", answer).strip() answer = STRUCTURED_LABEL_RE.sub(lambda match: label_replacement(match.group(1)), answer) if "почт" not in user_text.casefold(): answer = remove_email_fragments(answer) answer = voice_ready_address(answer, language) if language == "ru" and is_branch_name_only_query(user_text): answer = trim_branch_contact_details(answer) return WHITESPACE_RE.sub(" ", answer).strip() def humanize_assistant_answer( answer: str, user_text: str, language_code: str | None = None, ) -> str: language = normalize_language_code(language_code) or detect_conversation_language(user_text) or "ru" if language != "ru": return answer if not settings.humanizer_enabled or not settings.humanizer_fillers: return answer if not answer or FILLER_PREFIX_RE.match(answer): return answer if build_initial_details_ack(user_text, language) == answer: return answer if FILLER_BLOCK_RE.search(answer) or FILLER_BLOCK_RE.search(user_text): return answer if len(answer) < 25 or len(answer) > 220: return answer if random.random() > settings.humanizer_filler_rate: return answer filler = random.choice(settings.humanizer_fillers) return f"{filler}, {lowercase_initial_cyrillic(answer)}" def lowercase_initial_cyrillic(text: str) -> str: return CYRILLIC_TITLE_START_RE.sub( lambda match: f"{match.group(1).lower()}{match.group(2)}", text, count=1, ) def label_replacement(label: str) -> str: normalized = label.casefold() if normalized == "адрес": return "Он находится по адресу " if normalized == "телефон": return "Телефон " return "" def remove_email_fragments(text: str) -> str: text = re.sub( r"(?:электронная\s+почта|почта)\s*[:—-]?\s*\S+@\S+[.,;]?\s*", "", text, flags=re.IGNORECASE, ) return EMAIL_RE.sub("", text) def is_branch_name_only_query(user_text: str) -> bool: return bool(BRANCH_QUERY_RE.search(user_text)) and not CONTACT_REQUEST_RE.search(user_text) def trim_branch_contact_details(answer: str) -> str: short_answer = re.split( r"\s+(?:он\s+находится\s+по\s+(?:адресу|улице)|он\s+находится|" r"находится\s+по|адрес|телефон|почта)\b", answer, maxsplit=1, flags=re.IGNORECASE, )[0].strip(" .") sentence_match = re.search(r"[^.!?]*филиал[^.!?]*[.!?]?", short_answer, flags=re.IGNORECASE) if sentence_match: short_answer = sentence_match.group(0).strip(" .") if "филиал" not in short_answer.casefold(): return answer return f"{short_answer}. Если нужно, подскажу адрес или телефон." def voice_ready_address(text: str, language_code: str | None = None) -> str: language = normalize_language_code(language_code) or "ru" slash_word = "бөлшек" if language == "kk" else "дробь" text = ADDRESS_WITH_SLASH_RE.sub( lambda match: ( f"{number_to_words(int(match.group(1)), language)} {slash_word} " f"{number_to_words(int(match.group(2)), language)}" ), text, ) text = NUMBER_WITH_OPTIONAL_LETTER_RE.sub( lambda match: number_token_to_words(match.group(1), match.group(2), language), text, ) return WHITESPACE_RE.sub(" ", text).strip() def number_token_to_words(number_text: str, letter: str | None, language_code: str | None = None) -> str: words = number_to_words(int(number_text), language_code) if letter: return f"{words} {letter.upper()}" return words def number_to_words(number: int, language_code: str | None = None) -> str: if normalize_language_code(language_code) == "kk": return number_to_kazakh_words(number) return number_to_russian_words(number) def number_to_russian_words(number: int) -> str: if number < 10: return ONES[number] if number < 20: return TEENS[number] if number < 100: tens = number // 10 * 10 rest = number % 10 return f"{TENS[tens]} {ONES[rest]}" if rest else TENS[tens] hundreds = number // 100 * 100 rest = number % 100 if rest: return f"{HUNDREDS[hundreds]} {number_to_russian_words(rest)}" return HUNDREDS[hundreds] def number_to_kazakh_words(number: int) -> str: if number < 10: return KAZAKH_ONES[number] if number < 100: tens = number // 10 * 10 rest = number % 10 return f"{KAZAKH_TENS[tens]} {KAZAKH_ONES[rest]}" if rest else KAZAKH_TENS[tens] hundreds = number // 100 * 100 rest = number % 100 if rest: return f"{KAZAKH_HUNDREDS[hundreds]} {number_to_kazakh_words(rest)}" return KAZAKH_HUNDREDS[hundreds] def audio_format_label(output_format: str) -> str: if output_format.startswith("mp3_"): return "mp3" if output_format.startswith("pcm_"): return "pcm" if output_format.startswith("wav_"): return "wav" return output_format def audio_sample_rate(output_format: str) -> int: parts = output_format.split("_") if len(parts) >= 2 and parts[1].isdigit(): return int(parts[1]) return 16000 def audio_mime_type(output_format: str) -> str: if output_format.startswith("mp3_"): return "audio/mpeg" if output_format.startswith("wav_"): return "audio/wav" return "audio/pcm" def prepare_tts_text(text: str, language_code: str | None = None) -> str: return ( voice_ready_address(text, language_code) .replace("Айнур", "Ай-нур") .replace("айнур", "ай-нур") .replace("QazaqGasAimaq", "Казак Газ Аймак") ) @app.websocket("/ws") async def websocket_endpoint(websocket: WebSocket) -> None: session = VoiceSession(websocket) try: await session.run() except WebSocketDisconnect: session.closed = True except Exception as exc: session.closed = True try: await websocket.send_json({"type": "error", "message": str(exc)}) except Exception: pass