From d756003036b4a3bcdd4d82a3ae1b44ff2465476c Mon Sep 17 00:00:00 2001 From: Konturai DevOps Date: Fri, 1 May 2026 17:18:35 +0500 Subject: [PATCH] . --- .env.example | 4 +- core/kazakh_names.py | 591 ++++++++++++++++++++++++++++++++++++++++ core/session.py | 115 +++++++- crm_client.py | 114 ++++++++ docker-compose.yml | 8 + main.py | 4 + providers/base.py | 21 +- providers/llm.py | 28 +- providers/stt.py | 120 ++++++-- providers/stt_openai.py | 9 +- 10 files changed, 977 insertions(+), 37 deletions(-) create mode 100644 core/kazakh_names.py create mode 100644 crm_client.py diff --git a/.env.example b/.env.example index 3c92859..d07b6b3 100644 --- a/.env.example +++ b/.env.example @@ -9,7 +9,7 @@ LLM_PROVIDER=openai OPENAI_API_KEY= OPENAI_BASE_URL= OPENAI_LLM_MODEL=gpt-4o-mini -OPENAI_LLM_SYSTEM_PROMPT=You are Ainur, a concise Russian-speaking voice agent for the DigiOps contact center. Always answer in Russian, naturally, briefly, and in one calm consistent speaking style for voice playback. Never say waiting fillers like "минуту" or "секунду". If a short holding phrase is needed, the voice service handles it separately, so answer directly and do not add your own waiting phrase. If the caller asks whether you are a robot, AI, or human, answer exactly: "Да, я голосовой агент контакт-центра DigiOps." If the caller says goodbye, thanks, says nothing else is needed, or ends the conversation, respond briefly and politely without filler phrases. Read phone numbers digit by digit. Prefer clear stress-friendly wording and avoid ambiguous stress-heavy wording. Do not use markdown, URLs, or tables. +OPENAI_LLM_SYSTEM_PROMPT=Ты Айнур, лаконичный русскоязычный голосовой оператор контакт-центра DigiOps. Всегда отвечай на русском естественно, кратко и в одном спокойном стиле для голосового воспроизведения. Никогда не говори фразы ожидания вроде "минуту" или "секунду". Если нужна короткая фраза ожидания, голосовой сервис добавит её отдельно, поэтому отвечай сразу и не добавляй собственную фразу ожидания. Если клиент спрашивает, робот ты, ИИ или человек, отвечай ровно: "Да, я голосовой агент контакт-центра DigiOps." Не завершай диалог самостоятельно. Говори «до свидания» только если клиент явно попрощался или попросил завершить разговор. Простые «спасибо» или «благодарю» не считай просьбой завершить разговор. Телефонные номера читай по цифрам. Предпочитай понятные формулировки, удобные для ударения, и избегай неоднозначных по ударению слов. Не используй Markdown, URL или таблицы. OPENAI_TIMEOUT_SECONDS=30 OPENAI_STT_MODEL=whisper-1 OPENAI_STT_LANGUAGE=ru @@ -19,7 +19,7 @@ OLLAMA_BASE_URL=http://host.docker.internal:11434 OLLAMA_LLM_MODEL=qwen2.5:1.5b OLLAMA_TIMEOUT_SECONDS=20 OLLAMA_LLM_TEMPERATURE=0.3 -OLLAMA_LLM_SYSTEM_PROMPT=Ты Айнур, голосовой ИИ-оператор. Всегда отвечай только на русском, коротко и по делу. +OLLAMA_LLM_SYSTEM_PROMPT=Ты Айнур, голосовой ИИ-оператор. Всегда отвечай только на русском, коротко и по делу. Не завершай диалог самостоятельно и говори «до свидания» только если клиент явно попрощался или попросил завершить разговор. OLLAMA_LLM_MAX_CONTEXT_MESSAGES=4 OLLAMA_LLM_NUM_PREDICT=64 OLLAMA_LLM_NUM_CTX=1024 diff --git a/core/kazakh_names.py b/core/kazakh_names.py new file mode 100644 index 0000000..19e3fc3 --- /dev/null +++ b/core/kazakh_names.py @@ -0,0 +1,591 @@ +"""Kazakh first names for biasing the STT during the initial name-capture turn. + +Sized for ElevenLabs Scribe v2 *batch* keyterms cap (1000 entries, 50 chars +each). The Realtime path is intentionally skipped during the first turn — +see CallSession._start_live_stt_stream — so we use the larger batch ceiling +to maximize coverage for Kazakh-origin names. Russian-origin names are +excluded; they are already recognized correctly without biasing. +""" + +from __future__ import annotations + + +TOP_KAZAKH_NAMES: tuple[str, ...] = ( + # ---- Male ---- + "Абай", + "Абдрахман", + "Абдулла", + "Абдулмажид", + "Абдулжалил", + "Абжан", + "Абильхан", + "Абубакир", + "Абыз", + "Абылай", + "Абылайхан", + "Адал", + "Адахан", + "Адельжан", + "Адиль", + "Адильбек", + "Адильхан", + "Азамат", + "Азат", + "Азиз", + "Азизбек", + "Айбар", + "Айбат", + "Айбек", + "Айбол", + "Айболат", + "Айдан", + "Айдар", + "Айдарбек", + "Айдархан", + "Айдос", + "Айкын", + "Айнабек", + "Айнур", + "Айсен", + "Айсултан", + "Айткали", + "Айткеш", + "Айтуар", + "Акан", + "Акбар", + "Акжан", + "Акжол", + "Аким", + "Аксултан", + "Алау", + "Алдар", + "Алдияр", + "Алибек", + "Алим", + "Алимхан", + "Алишер", + "Алмагамбет", + "Алмаз", + "Алмас", + "Алпамыс", + "Алтай", + "Алтынбек", + "Аман", + "Аманбай", + "Аманжол", + "Аманкос", + "Аманкельды", + "Амангельды", + "Амир", + "Амиржан", + "Амирхан", + "Амре", + "Анар", + "Анарбай", + "Анарбек", + "Анвар", + "Аркат", + "Ардак", + "Арман", + "Арнур", + "Арсен", + "Арслан", + "Арыстан", + "Асан", + "Асет", + "Асим", + "Аскар", + "Аскат", + "Аскер", + "Аслан", + "Атабек", + "Аубакир", + "Ауез", + "Аян", + "Аяс", + "Бакдаулет", + "Бакир", + "Бактияр", + "Бактыбай", + "Бакыт", + "Бакытбек", + "Бакытжан", + "Балапан", + "Балгабай", + "Барлыбай", + "Бауыржан", + "Баян", + "Бекарыс", + "Бекасыл", + "Бекболат", + "Бекжан", + "Бекзат", + "Бексултан", + "Берген", + "Берик", + "Биржан", + "Болат", + "Болатбек", + "Болатхан", + "Бостандык", + "Бухар", + "Габидолла", + "Габит", + "Гали", + "Галым", + "Галымжан", + "Гани", + "Гумар", + "Гумарбек", + "Дамир", + "Данияр", + "Дархан", + "Дастан", + "Даулет", + "Дауренбек", + "Даурен", + "Дидар", + "Динмухамед", + "Динислам", + "Дияз", + "Досмухамед", + "Досжан", + "Досымбек", + "Дулат", + "Едил", + "Едыге", + "Ельжан", + "Еламан", + "Ерасыл", + "Ерболат", + "Ерзат", + "Ержан", + "Ерик", + "Еркебулан", + "Еркетай", + "Еркин", + "Ерлан", + "Ерлик", + "Ермагамбет", + "Ермек", + "Ермухан", + "Ернар", + "Ернур", + "Есенгали", + "Есим", + "Ескендир", + "Есламбек", + "Жайдар", + "Жакау", + "Жалгас", + "Жамал", + "Жамшидбек", + "Жанабат", + "Жанат", + "Жанболат", + "Жангильды", + "Жангир", + "Жанибек", + "Жасулан", + "Жаугашты", + "Жиен", + "Жолан", + "Жомарт", + "Жубан", + "Жумабай", + "Жумабек", + "Жумагали", + "Жумаш", + "Зейнолла", + "Зекен", + "Зейнулла", + "Зердебай", + "Идрис", + "Илиас", + "Ильяс", + "Иман", + "Имран", + "Иса", + "Искандер", + "Кабылбек", + "Кадыр", + "Кадыржан", + "Кайдар", + "Кайрат", + "Кайыргали", + "Кайсар", + "Канат", + "Канатбек", + "Канжар", + "Канипа", + "Карим", + "Касен", + "Касым", + "Кенес", + "Кенесары", + "Кенжебай", + "Куандык", + "Куаныш", + "Куат", + "Курманбай", + "Кылыш", + "Лукман", + "Магжан", + "Мадияр", + "Майлы", + "Майлыбай", + "Максат", + "Максут", + "Манар", + "Манас", + "Манатай", + "Маралбек", + "Марат", + "Мардан", + "Маулен", + "Мерген", + "Мирас", + "Мукагали", + "Мукаш", + "Муратбек", + "Мурат", + "Мухамед", + "Мухамеджан", + "Мухтар", + "Мырзабек", + "Мырзахан", + "Мырзахмет", + "Назар", + "Наурыз", + "Несип", + "Нияз", + "Нурадил", + "Нурбек", + "Нурбол", + "Нурдаулет", + "Нуржан", + "Нурлан", + "Нурлыбек", + "Нурмагамбет", + "Нурмурат", + "Нурсадык", + "Нурсат", + "Нурсултан", + "Нуртас", + "Нурым", + "Олжас", + "Омар", + "Ораз", + "Оразбай", + "Оразалы", + "Орак", + "Оралбек", + "Оралхан", + "Орынбай", + "Райымбек", + "Райыс", + "Рамазан", + "Расул", + "Рауан", + "Рахман", + "Рахмет", + "Раушан", + "Ринат", + "Рустам", + "Рустем", + "Сабит", + "Сабыр", + "Сагиман", + "Сагынтай", + "Сагынбай", + "Саду", + "Саин", + "Сайрам", + "Сакен", + "Самат", + "Самет", + "Санжар", + "Сапар", + "Саркен", + "Сасан", + "Сейтжан", + "Серик", + "Серикбай", + "Серикбек", + "Серикжан", + "Султан", + "Султанбек", + "Султанмурат", + "Сұңғат", + "Сұңқар", + "Талгат", + "Талгатбек", + "Тамерлан", + "Танат", + "Темирбек", + "Темирболат", + "Темирлан", + "Темиржан", + "Тилеген", + "Тимур", + "Толкын", + "Толеген", + "Тулеген", + "Турар", + "Турсын", + "Турсунбай", + "Узакбай", + "Улан", + "Уразалы", + "Усербай", + "Утеген", + "Файзулла", + "Фарит", + "Фархад", + "Хабиб", + "Хайрат", + "Хамза", + "Чингиз", + "Шалабай", + "Шамиль", + "Шерхан", + "Шернияз", + "Шынгысхан", + "Эльдар", + "Эрнар", + "Юсуф", + "Якуб", + # ---- Female ---- + "Аделя", + "Адина", + "Адия", + "Азиля", + "Айгерим", + "Айгуль", + "Айдай", + "Айдана", + "Айзада", + "Айжамал", + "Айжан", + "Айзере", + "Айман", + "Айна", + "Айнагуль", + "Айнура", + "Айсара", + "Айсауле", + "Айсулу", + "Айша", + "Айя", + "Айяна", + "Акбота", + "Акерке", + "Аксана", + "Акмарал", + "Акниет", + "Акшолпан", + "Алина", + "Алия", + "Алмагуль", + "Алтын", + "Алтынай", + "Алуа", + "Альбина", + "Альмира", + "Альфия", + "Анел", + "Анелия", + "Аннур", + "Аружан", + "Аруна", + "Асем", + "Асемгуль", + "Асель", + "Асима", + "Асия", + "Асылзада", + "Асылзат", + "Асылым", + "Аяжан", + "Аяна", + "Аяулым", + "Багдат", + "Багила", + "Бадигуль", + "Балауса", + "Балгуль", + "Балжан", + "Балнур", + "Балым", + "Бахытгуль", + "Бахытжан", + "Бибигуль", + "Бибинур", + "Бикеш", + "Венера", + "Газиза", + "Гайша", + "Галия", + "Гаухар", + "Гаухарбану", + "Гулим", + "Гулбану", + "Гулбаршын", + "Гулдана", + "Гулден", + "Гулжан", + "Гулзада", + "Гулзара", + "Гулмира", + "Гулназ", + "Гулнар", + "Гулсая", + "Гулсара", + "Гулсина", + "Гулсум", + "Гулшат", + "Гульдана", + "Гульден", + "Гульжан", + "Гульзада", + "Гульзара", + "Гульмира", + "Гульназ", + "Гульнар", + "Гульсая", + "Гульсина", + "Гульсум", + "Гульшат", + "Гульбану", + "Гульбаршын", + "Гулнара", + "Гультас", + "Дамира", + "Дана", + "Данара", + "Данель", + "Дария", + "Дарига", + "Дарика", + "Дилда", + "Дильдар", + "Дильнара", + "Дильназ", + "Дина", + "Динара", + "Динель", + "Жадыра", + "Жайна", + "Жанаргуль", + "Жанар", + "Жанбал", + "Жанна", + "Жания", + "Жасмин", + "Жибек", + "Жулдыз", + "Зайнаб", + "Замзагуль", + "Замира", + "Зара", + "Зарема", + "Зарина", + "Зарифа", + "Захира", + "Зейнеп", + "Зира", + "Зулайхан", + "Зухра", + "Иасмин", + "Инара", + "Индира", + "Камажай", + "Камила", + "Камилла", + "Карлыга", + "Карлыгаш", + "Касиет", + "Кулжайнар", + "Кулзада", + "Куляш", + "Куралай", + "Кымбат", + "Лаура", + "Лейла", + "Леяла", + "Лиана", + "Луиза", + "Мадина", + "Майра", + "Майя", + "Мақпал", + "Малика", + "Манара", + "Манагуль", + "Маржан", + "Маржанкуль", + "Мариямкуль", + "Мариям", + "Мариам", + "Меруерт", + "Мира", + "Молдир", + "Мунира", + "Назым", + "Назгуль", + "Назифа", + "Наргиз", + "Нурай", + "Нурбала", + "Нургуль", + "Нуржамал", + "Нурзада", + "Нуриля", + "Нурлыгуль", + "Перизат", + "Райхан", + "Рузия", + "Сабина", + "Сагила", + "Сайраш", + "Салима", + "Салтанат", + "Самал", + "Самира", + "Сандугаш", + "Сапура", + "Сара", + "Саулет", + "Сауле", + "Сахида", + "Симбат", + "Сулу", + "Сулушаш", + "Тазагуль", + "Тогжан", + "Толгана", + "Толганай", + "Томирис", + "Турсынай", + "Тұрсынгүл", + "Улжан", + "Улпан", + "Умит", + "Умитжан", + "Фариза", + "Фарзана", + "Фатима", + "Хадиша", + "Хайша", + "Шамшия", + "Шарбану", + "Шарифа", + "Шолпан", + "Шынар", + "Эльмира", + "Энлик", + "Юлдыз", +) + + +__all__ = ("TOP_KAZAKH_NAMES",) diff --git a/core/session.py b/core/session.py index 3bf713c..4ad36e4 100644 --- a/core/session.py +++ b/core/session.py @@ -15,6 +15,7 @@ from collections.abc import Iterable from enum import Enum from realtime_voice_service.core.filler_audio import FillerAudioLibrary +from realtime_voice_service.core.kazakh_names import TOP_KAZAKH_NAMES from realtime_voice_service.core.vad import BaseVAD, SileroVADDetector from realtime_voice_service.providers.base import BaseLLM, BaseSTT, BaseSTTStream, BaseTTS, MockLLM, MockSTT, MockTTS from realtime_voice_service.transports.base import BaseMediaTransport @@ -53,7 +54,7 @@ FILLER_AUDIO_DELAY_MS = _filler_delay_ms() TTS_CHUNK_SOFT_MIN_CHARS = _tts_chunk_soft_min_chars() TTS_CHUNK_SOFT_MIN_WORDS = _tts_chunk_soft_min_words() -_FAREWELL_MARKERS: frozenset[str] = frozenset({ +_COURTESY_OR_FAREWELL_MARKERS: frozenset[str] = frozenset({ "спасибо", "благодарю", "до свидания", @@ -77,6 +78,48 @@ _FAREWELL_MARKERS: frozenset[str] = frozenset({ "bye", }) +_TERMINAL_DIRECT_FAREWELL_MARKERS: frozenset[str] = frozenset({ + "до свидания", + "всего доброго", + "пока", + "до встречи", + "хорошего дня", + "goodbye", + "bye", +}) + +_TERMINAL_CLOSE_REQUEST_MARKERS: frozenset[str] = frozenset({ + "можно завершить", + "можем завершить", + "можете завершать", + "давайте завершим", + "давайте закончим", + "завершить разговор", + "закончить разговор", + "заканчиваем разговор", + "завершайте разговор", + "кладу трубку", + "положу трубку", +}) + +_TERMINAL_NO_MORE_HELP_MARKERS: frozenset[str] = frozenset({ + "ничего не нужно", + "ничего не надо", + "мне ничего не нужно", + "мне ничего не надо", + "больше ничего", + "больше ничего не нужно", + "больше ничего не надо", + "ничего больше не нужно", + "ничего больше не надо", + "это все", + "это всё", + "на этом все", + "на этом всё", + "вопросов нет", + "больше вопросов нет", +}) + _FAREWELL_QUERY_MARKERS: frozenset[str] = frozenset({ "вопрос", "подскаж", @@ -90,7 +133,6 @@ _FAREWELL_QUERY_MARKERS: frozenset[str] = frozenset({ "нужно ли", "нужно еще", "нужно ещё", - "хочу", "еще", "ещё", "теперь", @@ -107,22 +149,63 @@ _FAREWELL_QUERY_MARKERS: frozenset[str] = frozenset({ "нужна", "нужны", }) + +_FAREWELL_HARD_QUERY_MARKERS: frozenset[str] = frozenset( + marker for marker in _FAREWELL_QUERY_MARKERS if marker not in {"вопрос", "можно"} +) FAREWELL_RESPONSE_TEXT = "Спасибо за обращение в DigiOps. Всего доброго, до свидания!" +def _contains_voice_phrase(normalized_text: str, phrase: str) -> bool: + phrase_parts = [re.escape(part) for part in phrase.split()] + if not phrase_parts: + return False + phrase_pattern = r"\s+".join(phrase_parts) + return re.search( + rf"(? bool: + phrase_parts = [re.escape(part) for part in phrase.split()] + if not phrase_parts: + return False + phrase_pattern = r"\s+".join(phrase_parts) + return re.search( + rf"(? bool: normalized = _voice_text_key(text) - return any(marker in normalized for marker in _FAREWELL_MARKERS) + return any(_contains_voice_phrase(normalized, marker) for marker in _COURTESY_OR_FAREWELL_MARKERS) def _is_terminal_farewell(text: str) -> bool: normalized = _voice_text_key(text) - if not any(marker in normalized for marker in _FAREWELL_MARKERS): + has_direct_farewell = any( + _contains_voice_phrase(normalized, marker) + for marker in _TERMINAL_DIRECT_FAREWELL_MARKERS + ) + has_close_request = any( + _contains_voice_phrase(normalized, marker) + for marker in _TERMINAL_CLOSE_REQUEST_MARKERS + ) + has_no_more_help = any( + _ends_with_voice_phrase(normalized, marker) + for marker in _TERMINAL_NO_MORE_HELP_MARKERS + ) + if not (has_direct_farewell or has_close_request or has_no_more_help): return False - if any(marker in normalized for marker in _FAREWELL_QUERY_MARKERS): + query_markers = _FAREWELL_QUERY_MARKERS if has_direct_farewell else _FAREWELL_HARD_QUERY_MARKERS + if any(marker in normalized for marker in query_markers): return False word_count = len(re.findall(r"[^\W\d_]+(?:[-'][^\W\d_]+)*", normalized, flags=re.UNICODE)) - return word_count <= 10 + return word_count <= 16 SEMANTIC_CONTINUATION_TOKENS = { @@ -1405,7 +1488,12 @@ class CallSession: len(audio_bytes), _audio_duration_ms(audio_bytes, sample_rate_hz=self.transport.sample_rate_hz), ) - pending_transcript = (await self._stt.transcribe(audio_bytes)).strip() + pending_transcript = ( + await self._stt.transcribe( + audio_bytes, + keyterms=self._name_capture_keyterms() if self._awaiting_customer_name else None, + ) + ).strip() except Exception: LOGGER.exception("realtime session %s failed to transcribe pending unanswered audio", self.session_id) continue @@ -1421,6 +1509,9 @@ class CallSession: ) return pending_texts + def _name_capture_keyterms(self) -> tuple[str, ...]: + return TOP_KAZAKH_NAMES + def _build_name_collection_response(self, transcript: str) -> str | None: if not self._awaiting_customer_name: return None @@ -1880,7 +1971,10 @@ class CallSession: self._reset_semantic_endpointing() try: LOGGER.info("realtime session %s live STT stream start requested", self.session_id) - live_stt_stream = await self._stt.start_stream(partial_callback=self._handle_partial_transcript) + live_stt_stream = await self._stt.start_stream( + partial_callback=self._handle_partial_transcript, + keyterms=self._name_capture_keyterms() if self._awaiting_customer_name else None, + ) except Exception: LOGGER.exception("realtime session %s failed to start live STT stream", self.session_id) return @@ -1953,7 +2047,10 @@ class CallSession: len(utterance_audio), _audio_duration_ms(utterance_audio, sample_rate_hz=self.transport.sample_rate_hz), ) - return await self._stt.transcribe(utterance_audio) + return await self._stt.transcribe( + utterance_audio, + keyterms=self._name_capture_keyterms() if self._awaiting_customer_name else None, + ) def _maybe_dump_audio(self, audio_bytes: bytes) -> None: if str(os.getenv("ENABLE_AUDIO_DUMP", "")).strip().lower() not in {"1", "true", "yes", "on"}: diff --git a/crm_client.py b/crm_client.py new file mode 100644 index 0000000..7eac0b3 --- /dev/null +++ b/crm_client.py @@ -0,0 +1,114 @@ +from __future__ import annotations + +import base64 +import hashlib +import hmac +import json +import logging +import os +from datetime import datetime, timedelta, timezone +from typing import Any + +import httpx + +LOGGER = logging.getLogger("uvicorn.error") + + +def _b64url_encode(raw: bytes) -> str: + return base64.urlsafe_b64encode(raw).decode("utf-8").rstrip("=") + + +def _sign_hs256(message: str, secret: str) -> str: + digest = hmac.new( + secret.encode("utf-8"), + message.encode("utf-8"), + hashlib.sha256, + ).digest() + return _b64url_encode(digest) + + +def _issue_service_token(secret: str) -> str: + now = datetime.now(timezone.utc) + header = {"alg": "HS256", "typ": "JWT"} + payload = { + "sub": "svc:realtime-voice", + "username": "realtime-voice", + "role": "admin", + "auth_source": "service", + "iat": int(now.timestamp()), + "exp": int((now + timedelta(seconds=300)).timestamp()), + } + encoded_header = _b64url_encode(json.dumps(header, separators=(",", ":")).encode("utf-8")) + encoded_payload = _b64url_encode(json.dumps(payload, separators=(",", ":")).encode("utf-8")) + message = f"{encoded_header}.{encoded_payload}" + return f"{message}.{_sign_hs256(message, secret)}" + + +def _enabled() -> bool: + return os.getenv("CRM_INTERACTION_ENABLED", "0").strip().lower() in {"1", "true", "yes"} + + +def _base_url() -> str: + return str(os.getenv("CRM_INTERACTION_SERVICE_URL", "http://interaction-service:8000")).rstrip("/") + + +def _secret() -> str: + return str(os.getenv("CRM_APP_TOKEN_SECRET", "")).strip() + + +def _queue_id() -> str | None: + raw = os.getenv("CRM_QUEUE_ID", "").strip() + return raw or None + + +async def create_interaction(*, call_id: str) -> str | None: + if not _enabled(): + return None + secret = _secret() + if not secret: + LOGGER.warning("crm: CRM_APP_TOKEN_SECRET not configured, skipping interaction creation") + return None + token = _issue_service_token(secret) + payload: dict[str, Any] = { + "channel": "voice", + "subject": f"AI Voice Call [{call_id}]", + "priority": 3, + } + queue = _queue_id() + if queue: + payload["queue_id"] = queue + try: + async with httpx.AsyncClient(timeout=8.0) as client: + resp = await client.post( + f"{_base_url()}/interactions", + json=payload, + headers={"Authorization": f"Bearer {token}"}, + ) + resp.raise_for_status() + data = resp.json() + interaction_id = str(data.get("interaction_id") or "").strip() + LOGGER.info("crm: interaction created interaction_id=%s call_id=%s", interaction_id, call_id) + return interaction_id or None + except Exception: + LOGGER.exception("crm: failed to create interaction call_id=%s", call_id) + return None + + +async def close_interaction(interaction_id: str) -> None: + if not _enabled(): + return + secret = _secret() + if not secret: + return + token = _issue_service_token(secret) + try: + async with httpx.AsyncClient(timeout=8.0) as client: + resp = await client.patch( + f"{_base_url()}/interactions/{interaction_id}/status", + json={"status": "closed"}, + headers={"Authorization": f"Bearer {token}"}, + ) + resp.raise_for_status() + LOGGER.info("crm: interaction closed interaction_id=%s", interaction_id) + except Exception: + LOGGER.exception("crm: failed to close interaction interaction_id=%s", interaction_id) diff --git a/docker-compose.yml b/docker-compose.yml index dbbe4d4..56fc169 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -18,3 +18,11 @@ services: - "api.elevenlabs.io:34.8.184.191" - "host.docker.internal:host-gateway" restart: unless-stopped + networks: + - default + - call_center_net + +networks: + call_center_net: + name: call-center_default + external: true diff --git a/main.py b/main.py index 2b51814..0c36686 100644 --- a/main.py +++ b/main.py @@ -9,6 +9,7 @@ from contextlib import asynccontextmanager from fastapi import FastAPI, WebSocket +from realtime_voice_service import crm_client from realtime_voice_service.core.filler_audio import FillerAudioLibrary from realtime_voice_service.core.session import CallSession from realtime_voice_service.core.vad import SileroVADDetector @@ -200,6 +201,7 @@ class RealtimeVoiceService: async def _run_transport_session(self, transport: BaseMediaTransport) -> None: session = self._build_session(transport) + interaction_id = await crm_client.create_interaction(call_id=session.session_id) async with self._track_session(session): LOGGER.info( "starting realtime session %s via %s sample_rate=%s frame_ms=%s frame_bytes=%s", @@ -210,6 +212,8 @@ class RealtimeVoiceService: transport.frame_bytes, ) await session.run() + if interaction_id: + await crm_client.close_interaction(interaction_id) def _build_session(self, transport: BaseMediaTransport) -> CallSession: return CallSession( diff --git a/providers/base.py b/providers/base.py index ef228cc..02298ff 100644 --- a/providers/base.py +++ b/providers/base.py @@ -9,6 +9,7 @@ from collections.abc import AsyncGenerator from collections.abc import AsyncIterable from collections.abc import Awaitable from collections.abc import Callable +from collections.abc import Sequence from dataclasses import dataclass @@ -31,15 +32,21 @@ class BaseSTTStream(ABC): class BaseSTT(ABC): @abstractmethod - async def transcribe(self, audio_bytes: bytes) -> str: + async def transcribe( + self, + audio_bytes: bytes, + *, + keyterms: Sequence[str] | None = None, + ) -> str: raise NotImplementedError async def start_stream( self, *, partial_callback: PartialTranscriptCallback | None = None, + keyterms: Sequence[str] | None = None, ) -> BaseSTTStream | None: - del partial_callback + del partial_callback, keyterms return None @@ -76,7 +83,13 @@ class MockSTT(BaseSTT): self._scripted_transcripts = list(scripted_transcripts or []) self._call_count = 0 - async def transcribe(self, audio_bytes: bytes) -> str: + async def transcribe( + self, + audio_bytes: bytes, + *, + keyterms: Sequence[str] | None = None, + ) -> str: + del keyterms await asyncio.sleep(self._latency_ms / 1000.0) self._call_count += 1 if self._scripted_transcripts: @@ -88,7 +101,9 @@ class MockSTT(BaseSTT): self, *, partial_callback: PartialTranscriptCallback | None = None, + keyterms: Sequence[str] | None = None, ) -> BaseSTTStream | None: + del keyterms return _MockSTTStream(parent=self, partial_callback=partial_callback) diff --git a/providers/llm.py b/providers/llm.py index 87c942a..ca9bc39 100644 --- a/providers/llm.py +++ b/providers/llm.py @@ -15,6 +15,24 @@ from realtime_voice_service.providers.base import LLMStreamEvent LOGGER = logging.getLogger("uvicorn.error") +_CONVERSATION_CLOSE_POLICY = ( + "Правило завершения диалога: ИИ-оператор не должен самостоятельно завершать разговор. " + "Не говори «до свидания», «всего доброго», «хорошего дня» и другие финальные прощания, " + "если клиент явно не попрощался или не попросил завершить разговор. " + "Обычная благодарность вроде «спасибо» или «благодарю» не является просьбой завершить разговор; " + "на нее отвечай коротко и оставляй диалог открытым." +) +_CONVERSATION_CLOSE_POLICY_MARKER = "ИИ-оператор не должен самостоятельно завершать разговор" + + +def _with_conversation_close_policy(system_prompt: str) -> str: + normalized_prompt = str(system_prompt or "").strip() + if _CONVERSATION_CLOSE_POLICY_MARKER in normalized_prompt: + return normalized_prompt + if not normalized_prompt: + return _CONVERSATION_CLOSE_POLICY + return f"{normalized_prompt}\n\n{_CONVERSATION_CLOSE_POLICY}" + def _preview_text(text: str, *, limit: int = 160) -> str: normalized = " ".join(str(text or "").split()) @@ -81,14 +99,15 @@ class OpenAILLM(BaseLLM): self._api_key = str(api_key if api_key is not None else os.getenv("OPENAI_API_KEY", "")).strip() self._model = str(model or os.getenv("OPENAI_LLM_MODEL", "gpt-4o-mini")).strip() or "gpt-4o-mini" self._base_url = str(base_url or os.getenv("OPENAI_BASE_URL", "")).strip() or None - self._system_prompt = str( + raw_system_prompt = ( system_prompt if system_prompt is not None else os.getenv( "OPENAI_LLM_SYSTEM_PROMPT", "Ты — дружелюбный, живой и эмпатичный голосовой ИИ-ассистент. Отвечай кратко, как в реальном диалоге. Используй разговорный стиль. Чтобы синтезатор речи (TTS) читал аббревиатуры и английские термины без акцента, пиши их русскими буквами так, как они произносятся (например, 'ай-ти' вместо 'IT', 'би-ту-би' вместо 'B2B', 'си-эр-эм' вместо 'CRM').", ) - ).strip() + ) + self._system_prompt = _with_conversation_close_policy(str(raw_system_prompt).strip()) self._timeout_seconds = max(float(timeout_seconds if timeout_seconds is not None else _timeout_seconds()), 1.0) self._temperature = max(min(float(temperature), 2.0), 0.0) self._reasoning_effort = str( @@ -647,7 +666,7 @@ class OllamaLLM(BaseLLM): str(base_url or os.getenv("OLLAMA_BASE_URL", "http://host.docker.internal:11434")).strip().rstrip("/") or "http://host.docker.internal:11434" ) - self._system_prompt = str( + raw_system_prompt = ( system_prompt if system_prompt is not None else os.getenv( @@ -657,7 +676,8 @@ class OllamaLLM(BaseLLM): "Ты — дружелюбный, живой и эмпатичный голосовой ИИ-ассистент. Отвечай кратко, как в реальном диалоге. Используй разговорный стиль. Чтобы синтезатор речи (TTS) читал аббревиатуры и английские термины без акцента, пиши их русскими буквами так, как они произносятся (например, 'ай-ти' вместо 'IT', 'би-ту-би' вместо 'B2B', 'си-эр-эм' вместо 'CRM').", ), ) - ).strip() + ) + self._system_prompt = _with_conversation_close_policy(str(raw_system_prompt).strip()) self._timeout_seconds = max( float(timeout_seconds if timeout_seconds is not None else self._read_float_env("OLLAMA_TIMEOUT_SECONDS", _timeout_seconds())), 1.0, diff --git a/providers/stt.py b/providers/stt.py index 760771f..3caf690 100644 --- a/providers/stt.py +++ b/providers/stt.py @@ -10,6 +10,7 @@ import logging import os import time import wave +from collections.abc import Sequence from urllib.parse import urlencode from realtime_voice_service.providers.base import BaseSTT @@ -22,6 +23,34 @@ LOGGER = logging.getLogger("uvicorn.error") _STREAM_COMMIT = object() _STREAM_CANCEL = object() +_REALTIME_MAX_KEYTERMS = 50 +_REALTIME_MAX_KEYTERM_CHARS = 20 +_BATCH_MAX_KEYTERMS = 1000 +_BATCH_MAX_KEYTERM_CHARS = 50 + + +def _normalize_keyterms( + keyterms: Sequence[str] | None, + *, + max_count: int, + max_chars: int, +) -> list[str]: + if not keyterms: + return [] + seen: set[str] = set() + normalized: list[str] = [] + for raw in keyterms: + term = " ".join(str(raw or "").split()) + if not term or len(term) > max_chars: + continue + if term in seen: + continue + seen.add(term) + normalized.append(term) + if len(normalized) >= max_count: + break + return normalized + def _api_base() -> str: return (os.getenv("ELEVENLABS_API_BASE", "https://api.elevenlabs.io").strip() or "https://api.elevenlabs.io").rstrip("/") @@ -528,7 +557,12 @@ class ElevenLabsSTT(BaseSTT): self._batch_min_audio_ms, ) - async def transcribe(self, audio_bytes: bytes) -> str: + async def transcribe( + self, + audio_bytes: bytes, + *, + keyterms: Sequence[str] | None = None, + ) -> str: if not audio_bytes: return "" if not self._api_key: @@ -537,12 +571,13 @@ class ElevenLabsSTT(BaseSTT): pcm_bytes, sample_rate_hz = self._extract_pcm(audio_bytes) LOGGER.info( "ElevenLabs STT transcribe start: input_bytes=%s extracted_pcm_bytes=%s sample_rate=%s duration_ms=%s " - "use_realtime=%s", + "use_realtime=%s keyterms=%s", len(audio_bytes), len(pcm_bytes), sample_rate_hz, _pcm_duration_ms(pcm_bytes, sample_rate_hz=sample_rate_hz), self._use_realtime, + len(keyterms) if keyterms else 0, ) if sample_rate_hz != self._target_sample_rate_hz: before_rate_hz = sample_rate_hz @@ -566,6 +601,7 @@ class ElevenLabsSTT(BaseSTT): transcript = await self._transcribe_realtime( pcm_bytes=pcm_bytes, sample_rate_hz=sample_rate_hz, + keyterms=keyterms, ) LOGGER.info( "ElevenLabs STT realtime transcript result: chars=%s transcript=%r", @@ -581,6 +617,7 @@ class ElevenLabsSTT(BaseSTT): transcript = await self._transcribe_batch( pcm_bytes=pcm_bytes, sample_rate_hz=sample_rate_hz, + keyterms=keyterms, ) LOGGER.info( "ElevenLabs STT batch transcript result: chars=%s transcript=%r", @@ -593,18 +630,30 @@ class ElevenLabsSTT(BaseSTT): self, *, partial_callback: PartialTranscriptCallback | None = None, + keyterms: Sequence[str] | None = None, ) -> BaseSTTStream | None: if not self._use_realtime: return None if not self._api_key: raise RuntimeError("ELEVENLABS_API_KEY is required for ElevenLabs STT") - websocket_url = self._build_realtime_websocket_url(sample_rate_hz=self._target_sample_rate_hz) + websocket_url = self._build_realtime_websocket_url( + sample_rate_hz=self._target_sample_rate_hz, + keyterms=keyterms, + ) + normalized_keyterm_count = len( + _normalize_keyterms( + keyterms, + max_count=_REALTIME_MAX_KEYTERMS, + max_chars=_REALTIME_MAX_KEYTERM_CHARS, + ) + ) LOGGER.info( - "ElevenLabs STT live stream starting: realtime_model=%s sample_rate=%s language=%s", + "ElevenLabs STT live stream starting: realtime_model=%s sample_rate=%s language=%s keyterms=%s", self._realtime_model_id, self._target_sample_rate_hz, self._language_code, + normalized_keyterm_count, ) stream = ElevenLabsRealtimeSTTStream( api_key=self._api_key, @@ -628,6 +677,7 @@ class ElevenLabsSTT(BaseSTT): *, pcm_bytes: bytes, sample_rate_hz: int, + keyterms: Sequence[str] | None = None, ) -> str: try: from websockets.exceptions import WebSocketException @@ -635,7 +685,10 @@ class ElevenLabsSTT(BaseSTT): except Exception as exc: # noqa: BLE001 raise RuntimeError("The `websockets` package is required for ElevenLabs realtime STT") from exc - websocket_url = self._build_realtime_websocket_url(sample_rate_hz=sample_rate_hz) + websocket_url = self._build_realtime_websocket_url( + sample_rate_hz=sample_rate_hz, + keyterms=keyterms, + ) chunk_bytes = max(int(sample_rate_hz * self._realtime_chunk_duration_ms / 1000.0) * 2, 320) started_monotonic = time.perf_counter() sent_chunks = 0 @@ -736,6 +789,7 @@ class ElevenLabsSTT(BaseSTT): *, pcm_bytes: bytes, sample_rate_hz: int, + keyterms: Sequence[str] | None = None, ) -> str: # TODO: Architectural Bottleneck: Рассмотреть замену STT на Deepgram WebSocket API для достижения true-streaming latency. try: @@ -749,16 +803,22 @@ class ElevenLabsSTT(BaseSTT): min_duration_ms=self._batch_min_audio_ms, ) wav_bytes = _pcm16le_to_wav_bytes(pcm_bytes, sample_rate_hz=sample_rate_hz) + normalized_keyterms = _normalize_keyterms( + keyterms, + max_count=_BATCH_MAX_KEYTERMS, + max_chars=_BATCH_MAX_KEYTERM_CHARS, + ) started_monotonic = time.perf_counter() LOGGER.info( "ElevenLabs batch STT request start: model=%s sample_rate=%s pcm_bytes=%s wav_bytes=%s " - "duration_ms=%s language=%s", + "duration_ms=%s language=%s keyterms=%s", self._model_id, sample_rate_hz, len(pcm_bytes), len(wav_bytes), _pcm_duration_ms(pcm_bytes, sample_rate_hz=sample_rate_hz), self._language_code, + len(normalized_keyterms), ) form = aiohttp.FormData() form.add_field("model_id", self._model_id) @@ -766,6 +826,8 @@ class ElevenLabsSTT(BaseSTT): form.add_field("diarize", "false") if self._language_code: form.add_field("language_code", self._language_code) + if normalized_keyterms: + form.add_field("keyterms", json.dumps(normalized_keyterms, ensure_ascii=False)) form.add_field( "file", wav_bytes, @@ -860,15 +922,26 @@ class ElevenLabsSTT(BaseSTT): return f"ElevenLabs realtime STT {message_type}: {detail}" return f"ElevenLabs realtime STT {message_type}" - def _build_realtime_websocket_url(self, *, sample_rate_hz: int) -> str: - query = { - "model_id": self._realtime_model_id, - "audio_format": f"pcm_{sample_rate_hz}", - "commit_strategy": "manual", - "include_timestamps": "false", - } + def _build_realtime_websocket_url( + self, + *, + sample_rate_hz: int, + keyterms: Sequence[str] | None = None, + ) -> str: + query: list[tuple[str, str]] = [ + ("model_id", self._realtime_model_id), + ("audio_format", f"pcm_{sample_rate_hz}"), + ("commit_strategy", "manual"), + ("include_timestamps", "false"), + ] if self._language_code: - query["language_code"] = self._language_code + query.append(("language_code", self._language_code)) + for term in _normalize_keyterms( + keyterms, + max_count=_REALTIME_MAX_KEYTERMS, + max_chars=_REALTIME_MAX_KEYTERM_CHARS, + ): + query.append(("keyterms", term)) return f"{self._ws_base}/v1/speech-to-text/realtime?{urlencode(query)}" def _extract_pcm(self, audio_bytes: bytes) -> tuple[bytes, int]: @@ -926,7 +999,13 @@ class YandexSpeechKitSTT(BaseSTT): "iam" if self._iam_token else ("api-key" if self._api_key else "missing"), ) - async def transcribe(self, audio_bytes: bytes) -> str: + async def transcribe( + self, + audio_bytes: bytes, + *, + keyterms: Sequence[str] | None = None, + ) -> str: + del keyterms if not audio_bytes: return "" if not self._api_key and not self._iam_token: @@ -1061,16 +1140,21 @@ class FallbackSTT(BaseSTT): type(fallback).__name__, ) - async def transcribe(self, audio_bytes: bytes) -> str: + async def transcribe( + self, + audio_bytes: bytes, + *, + keyterms: Sequence[str] | None = None, + ) -> str: try: - return await self._primary.transcribe(audio_bytes) + return await self._primary.transcribe(audio_bytes, keyterms=keyterms) except Exception: LOGGER.exception( "Primary STT provider failed; falling back: primary=%s fallback=%s", type(self._primary).__name__, type(self._fallback).__name__, ) - return await self._fallback.transcribe(audio_bytes) + return await self._fallback.transcribe(audio_bytes, keyterms=keyterms) async def close(self) -> None: for provider in (self._primary, self._fallback): diff --git a/providers/stt_openai.py b/providers/stt_openai.py index 76ec087..1ba8064 100644 --- a/providers/stt_openai.py +++ b/providers/stt_openai.py @@ -6,6 +6,7 @@ import logging import os import time import wave +from collections.abc import Sequence from typing import Any from realtime_voice_service.providers.base import BaseSTT @@ -94,7 +95,13 @@ class OpenAISTT(BaseSTT): self._timeout_seconds, ) - async def transcribe(self, audio_bytes: bytes) -> str: + async def transcribe( + self, + audio_bytes: bytes, + *, + keyterms: Sequence[str] | None = None, + ) -> str: + del keyterms if not audio_bytes: return "" if not self._api_key: