Compare commits
27
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c261d08036 | ||
|
|
78dfe6a9e1 | ||
|
|
f978702a6c | ||
|
|
f82f27c035 | ||
|
|
75d105f636 | ||
|
|
9bf367abf4 | ||
|
|
3c5c233071 | ||
|
|
99d169ec67 | ||
|
|
9fdaa9472f | ||
|
|
1dc37d2764 | ||
|
|
8382dfa9ba | ||
|
|
8ced7a59e3 | ||
|
|
d2438b6954 | ||
|
|
3826f5704a | ||
|
|
2cb358e80d | ||
|
|
2f4a9795b5 | ||
|
|
48d1fabba1 | ||
|
|
66652845d8 | ||
|
|
1dcfaf46cf | ||
|
|
010a8dcab6 | ||
|
|
13ba8f8b56 | ||
|
|
9d918af12a | ||
|
|
ff85fa27fb | ||
|
|
82c89eff5a | ||
|
|
36c2acb011 | ||
|
|
fe9b3f3a80 | ||
|
|
0a829d26c7 |
@@ -0,0 +1,13 @@
|
||||
{
|
||||
"$schema": "https://json-schema.org/draft/2020-12/schema",
|
||||
"title": "AgentConnected",
|
||||
"type": "object",
|
||||
"required": ["event", "call_id", "escalation_id", "agent_id", "created_at"],
|
||||
"properties": {
|
||||
"event": { "const": "AgentConnected" },
|
||||
"call_id": { "type": "string" },
|
||||
"escalation_id": { "type": "string" },
|
||||
"agent_id": { "type": "string" },
|
||||
"created_at": { "type": "string", "format": "date-time" }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"$schema": "https://json-schema.org/draft/2020-12/schema",
|
||||
"title": "AgentNoAnswer",
|
||||
"type": "object",
|
||||
"required": ["event", "call_id", "escalation_id", "agent_id", "dial_outcome", "created_at"],
|
||||
"properties": {
|
||||
"event": { "const": "AgentNoAnswer" },
|
||||
"call_id": { "type": "string" },
|
||||
"escalation_id": { "type": "string" },
|
||||
"agent_id": { "type": ["string", "null"] },
|
||||
"dial_outcome": { "type": "string" },
|
||||
"attempt_count": { "type": "integer" },
|
||||
"created_at": { "type": "string", "format": "date-time" }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,18 @@
|
||||
{
|
||||
"$schema": "https://json-schema.org/draft/2020-12/schema",
|
||||
"title": "AgentReserved",
|
||||
"type": "object",
|
||||
"required": ["event", "call_id", "escalation_id", "agent_id", "created_at"],
|
||||
"properties": {
|
||||
"event": { "const": "AgentReserved" },
|
||||
"call_id": { "type": "string" },
|
||||
"escalation_id": { "type": "string" },
|
||||
"agent_id": { "type": "string" },
|
||||
"tenant_id": { "type": ["string", "null"] },
|
||||
"from_level": { "type": "string" },
|
||||
"to_level": { "type": "string" },
|
||||
"reason_code": { "type": "string" },
|
||||
"attempt_count": { "type": "integer" },
|
||||
"created_at": { "type": "string", "format": "date-time" }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
{
|
||||
"$schema": "https://json-schema.org/draft/2020-12/schema",
|
||||
"title": "AgentRinging",
|
||||
"type": "object",
|
||||
"required": ["event", "call_id", "escalation_id", "agent_id", "created_at"],
|
||||
"properties": {
|
||||
"event": { "const": "AgentRinging" },
|
||||
"call_id": { "type": "string" },
|
||||
"escalation_id": { "type": "string" },
|
||||
"agent_id": { "type": "string" },
|
||||
"created_at": { "type": "string", "format": "date-time" }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
{
|
||||
"$schema": "https://json-schema.org/draft/2020-12/schema",
|
||||
"title": "TransferCompleted",
|
||||
"type": "object",
|
||||
"required": ["event", "call_id", "escalation_id", "agent_id", "created_at"],
|
||||
"properties": {
|
||||
"event": { "const": "TransferCompleted" },
|
||||
"call_id": { "type": "string" },
|
||||
"escalation_id": { "type": "string" },
|
||||
"agent_id": { "type": "string" },
|
||||
"created_at": { "type": "string", "format": "date-time" }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"$schema": "https://json-schema.org/draft/2020-12/schema",
|
||||
"title": "TransferFailed",
|
||||
"type": "object",
|
||||
"required": ["event", "call_id", "created_at"],
|
||||
"properties": {
|
||||
"event": { "const": "TransferFailed" },
|
||||
"call_id": { "type": "string" },
|
||||
"escalation_id": { "type": ["string", "null"] },
|
||||
"agent_id": { "type": ["string", "null"] },
|
||||
"reason": { "type": ["string", "null"] },
|
||||
"error": { "type": ["string", "null"] },
|
||||
"created_at": { "type": "string", "format": "date-time" }
|
||||
}
|
||||
}
|
||||
@@ -10,10 +10,10 @@ ALLOW_LEGACY_HEADER_AUTH=0
|
||||
|
||||
AI_PROVIDER=openai_compatible
|
||||
AI_API_BASE=https://api.openai.com/v1
|
||||
AI_API_KEY=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
|
||||
AI_MODEL=gpt-4o-mini
|
||||
AI_API_KEY=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
|
||||
AI_MODEL=gpt-5-mini
|
||||
AI_TIMEOUT_SECONDS=30
|
||||
AI_VOICE_AI_TIMEOUT_SECONDS=10
|
||||
AI_VOICE_AI_TIMEOUT_SECONDS=15
|
||||
AI_WEB_SEARCH_ENABLED=1
|
||||
AI_WEB_SEARCH_MAX_RESULTS=5
|
||||
AI_WEB_SEARCH_GL=kz
|
||||
@@ -84,9 +84,9 @@ AI_VOICE_TTS_CACHE_DIR=/app/.data_local/ai_voice_tts_cache
|
||||
AI_VOICE_TTS_ELEVENLABS_API_KEY=sk_1c50faab05df13fd0ecde9aeca4abadcc7269ffdc15a03ad
|
||||
AI_VOICE_TTS_ELEVENLABS_API_BASE=https://api.elevenlabs.io
|
||||
AI_VOICE_TTS_ELEVENLABS_MODEL_ID=eleven_turbo_v2_5
|
||||
AI_VOICE_TTS_ELEVENLABS_RU_VOICE_ID=0ArNnoIAWKlT4WweaVMY
|
||||
AI_VOICE_TTS_ELEVENLABS_RU_VOICE_ID=ut2XM2wJyIZLTtW6lFzZ
|
||||
AI_VOICE_TTS_ELEVENLABS_RU_LANGUAGE_CODE=ru
|
||||
AI_VOICE_TTS_ELEVENLABS_KK_VOICE_ID=0ArNnoIAWKlT4WweaVMY
|
||||
AI_VOICE_TTS_ELEVENLABS_KK_VOICE_ID=ut2XM2wJyIZLTtW6lFzZ
|
||||
AI_VOICE_TTS_ELEVENLABS_KK_LANGUAGE_CODE=kk
|
||||
AI_VOICE_TTS_ELEVENLABS_OUTPUT_FORMAT=pcm_16000
|
||||
AI_VOICE_TTS_YANDEX_API_KEY=AQWJUMiUaXmbegxN4kgvM2XIlNqAPoBR5Wtq-40
|
||||
@@ -185,3 +185,10 @@ EVENT_BUS_EXCHANGE=mvpcc.domain.events
|
||||
ERROR_TELEGRAM_ALERTS=0
|
||||
ERROR_TELEGRAM_ALERT_BOT_TOKEN=8734026216:AAEK-xdgVR5al6rfKp_DIb-HSJZXGUG2Tvg
|
||||
ERROR_TELEGRAM_ALERT_CHANNEL_ID=t.me/konturaitelecom
|
||||
AI_OPERATOR_AGENT_NAME=Жанна
|
||||
AI_OPERATOR_COMPANY_NAME=Казакгаз Аймак
|
||||
AI_OPERATOR_VOICE_GREETING_RU=Добрый день! Меня зовут Жанна, я ИИ-оператор контакт-центра Казакгаз Аймак. Вам удобнее получить консультацию на казахском или на русском языке?
|
||||
AI_OPERATOR_VOICE_GREETING_KZ=Қайырлы күн! Менің атым Жанна, мен «Қазақгаз Аймақ» байланыс орталығында жасанды интеллект операторымын. Кеңесті қазақ немесе орыс тілінде алғыңыз келе ме?
|
||||
AI_OPERATOR_IDENTITY_REPLY_RU=Я Жанна, ИИ-оператор контакт-центра Казакгаз Аймак. Чем могу помочь?
|
||||
AI_OPERATOR_IDENTITY_REPLY_KZ=Мен Жаннамын, «Қазақгаз Аймақ» байланыс орталығының жасанды интеллект операторымын. Қалай көмектесе аламын?
|
||||
AI_OPERATOR_BASE_SYSTEM_PROMPT=Ты Жанна, единый ИИ-оператор контакт-центра Казакгаз Аймак для звонков, Telegram и других каналов. Всегда сохраняй одну и ту же личность: тебя зовут Жанна. Если клиент спрашивает, кто ты или как тебя зовут, отвечай, что ты Жанна. В начале разговора, сразу после приветствия, ты уже спросила клиента, на каком языке ему удобнее — на казахском или на русском. Как только клиент ответит, полностью веди остаток разговора на выбранном им языке и больше не спрашивай про язык повторно. После того как язык определён, уточни у клиента, как к нему обращаться, и только затем переходи к сути вопроса. Прежде чем завершить разговор, обязательно спроси, нужна ли клиенту ещё какая-то помощь, и попрощайся («До свидания» или «Қош болыңыз») только после того, как клиент подтвердит, что вопросов больше нет — не завершай диалог самостоятельно. Отвечай естественно, кратко и по делу. Когда говоришь о себе, используй женский род: могла, смогла, сделала, готова, проверила, нашла. Телефонные номера читай по цифрам. Не используй Markdown, URL или таблицы. DOMAIN SCOPE: ты отвечаешь ТОЛЬКО на вопросы, связанные с услугами газоснабжения: оплата газа, счётчики, ваучеры, технические условия, отключение газа, мобильное приложение, квитанции, передача показаний, контакты филиалов, безопасность газа. Если вопрос НЕ связан с газом (погода, спорт, политика, другие компании, личные советы, кулинария, рецепты, анекдоты, новости, фильмы), вежливо откажи БЕЗ поиска в базе знаний и БЕЗ попытки ответить: RU «Я могу ответить только на вопросы, связанные с услугами газоснабжения. Пожалуйста, задайте вопрос по теме газа.» KK «Мен тек газмен қамтамасыз ету қызметтеріне байланысты сұрақтарға жауап беремін. Газ тақырыбында сұрақ қойыңызшы.» Не пытайся ответить на вопрос не по теме, даже если знаешь ответ.
|
||||
|
||||
@@ -0,0 +1,4 @@
|
||||
ALTER TABLE escalations ADD COLUMN IF NOT EXISTS attempt_count INTEGER NOT NULL DEFAULT 0;
|
||||
ALTER TABLE escalations ADD COLUMN IF NOT EXISTS real_agent_id TEXT;
|
||||
ALTER TABLE escalations ADD COLUMN IF NOT EXISTS attempted_agent_ids_json TEXT NOT NULL DEFAULT '[]';
|
||||
CREATE INDEX IF NOT EXISTS ix_escalations_real_agent_id ON escalations(real_agent_id);
|
||||
@@ -0,0 +1,4 @@
|
||||
ALTER TABLE escalations ADD COLUMN attempt_count INTEGER NOT NULL DEFAULT 0;
|
||||
ALTER TABLE escalations ADD COLUMN real_agent_id TEXT;
|
||||
ALTER TABLE escalations ADD COLUMN attempted_agent_ids_json TEXT NOT NULL DEFAULT '[]';
|
||||
CREATE INDEX IF NOT EXISTS ix_escalations_real_agent_id ON escalations(real_agent_id);
|
||||
@@ -0,0 +1,2 @@
|
||||
ALTER TABLE kb_articles ADD COLUMN IF NOT EXISTS intent_code TEXT;
|
||||
CREATE INDEX IF NOT EXISTS idx_kb_articles_intent_code ON kb_articles(intent_code);
|
||||
@@ -0,0 +1,2 @@
|
||||
ALTER TABLE kb_articles ADD COLUMN intent_code TEXT;
|
||||
CREATE INDEX IF NOT EXISTS idx_kb_articles_intent_code ON kb_articles(intent_code);
|
||||
@@ -9,7 +9,7 @@ import os
|
||||
import re
|
||||
import time
|
||||
from threading import Lock
|
||||
from typing import Any
|
||||
from typing import Any, Iterable
|
||||
|
||||
import httpx
|
||||
from fastapi import Depends, FastAPI, HTTPException, Query
|
||||
@@ -23,6 +23,7 @@ from services.shared.ai_context_summary import (
|
||||
update_context_summary_from_assistant_turn,
|
||||
update_context_summary_from_user_turn,
|
||||
)
|
||||
from services.shared.intents import normalize_intent
|
||||
from services.shared.kb_localization import normalize_kb_language
|
||||
from services.shared.kb_search import search_kb_rows
|
||||
from services.shared.models import (
|
||||
@@ -2315,6 +2316,7 @@ def _openai_prompt(
|
||||
"article_id": article.article_id,
|
||||
"title": article.title,
|
||||
"snippet": _article_snippet(article),
|
||||
"intent_code": getattr(article, "intent_code", None),
|
||||
}
|
||||
for article in kb_results
|
||||
]
|
||||
@@ -2355,13 +2357,18 @@ def _request_structured_model_decision(
|
||||
if not _ai_api_base() or not _ai_api_key():
|
||||
raise RuntimeError("AI_API_BASE / AI_API_KEY are required for openai_compatible provider")
|
||||
started = time.perf_counter()
|
||||
payload = {
|
||||
"model": _ai_model(),
|
||||
"temperature": 0.2,
|
||||
"max_tokens": _ai_decision_max_tokens(),
|
||||
model = _ai_model()
|
||||
payload: dict[str, Any] = {
|
||||
"model": model,
|
||||
"response_format": {"type": "json_object"},
|
||||
"messages": messages,
|
||||
}
|
||||
if model.startswith("gpt-5"):
|
||||
payload["max_completion_tokens"] = _ai_decision_max_tokens()
|
||||
payload["reasoning_effort"] = "minimal"
|
||||
else:
|
||||
payload["temperature"] = 0.2
|
||||
payload["max_tokens"] = _ai_decision_max_tokens()
|
||||
effective_timeout = timeout_seconds if timeout_seconds is not None else _ai_timeout_seconds()
|
||||
with httpx.Client(timeout=effective_timeout) as client:
|
||||
response = client.post(
|
||||
@@ -2410,10 +2417,15 @@ def _openai_compatible_decision(
|
||||
)
|
||||
|
||||
|
||||
def _sanitize_decision(raw: dict[str, Any], *, fallback_language: str) -> dict[str, Any]:
|
||||
def _sanitize_decision(
|
||||
raw: dict[str, Any],
|
||||
*,
|
||||
fallback_language: str,
|
||||
known_topic_codes: Iterable[str] = (),
|
||||
) -> dict[str, Any]:
|
||||
decision = {
|
||||
"language": str(raw.get("language") or fallback_language or "ru"),
|
||||
"intent": str(raw.get("intent") or "unknown"),
|
||||
"intent": normalize_intent(raw.get("intent"), known_topic_codes=known_topic_codes),
|
||||
"reply_text": str(raw.get("reply_text") or "").strip(),
|
||||
"extracted_name": str(raw.get("extracted_name") or "").strip() or None,
|
||||
"confidence": float(raw.get("confidence") or 0.0),
|
||||
@@ -2611,7 +2623,13 @@ def _decide_reply(
|
||||
raw["_model"] = _ai_model()
|
||||
raw["_latency_ms"] = 1
|
||||
raw["_finish_reason"] = "stop"
|
||||
return _sanitize_decision(raw, fallback_language=language)
|
||||
return _sanitize_decision(
|
||||
raw,
|
||||
fallback_language=language,
|
||||
known_topic_codes=[
|
||||
code for article in kb_results if (code := getattr(article, "intent_code", None))
|
||||
],
|
||||
)
|
||||
|
||||
|
||||
def _update_ai_session_context_summary_from_user_turn(
|
||||
|
||||
@@ -135,8 +135,12 @@ def human_fallback_reply(language: str, *, is_greeting: bool = False) -> str:
|
||||
def operator_system_prompt(*, language: str, channel_label: str, is_voice: bool, config: Any | None = None) -> str:
|
||||
preferred_language = "Kazakh" if str(language or "").strip().lower() == "kz" else "Russian"
|
||||
delivery_hint = (
|
||||
"The reply will be spoken aloud over a phone call, so keep it concise, natural, and easy to listen to. "
|
||||
"Prefer one or two short sentences and at most one clarifying question. "
|
||||
"The reply will be spoken aloud over a phone call, so keep it concise, natural, and easy to listen to, "
|
||||
"but never at the cost of dropping a required fact. "
|
||||
"Use as many short sentences as needed to cover every material fact from the grounding kb_results "
|
||||
"snippet completely - required steps, codes, commands, deadlines, amounts, and conditions - typically "
|
||||
"two to four short sentences; never silently omit or shorten out a required step just to sound brief. "
|
||||
"At most one clarifying question per turn. "
|
||||
"The text-to-speech engine reads exactly what you write, digit by digit, with no number formatting of its own, "
|
||||
"so never output bare digits for phone numbers, hotline numbers, or dates — always spell them out in words "
|
||||
"the way a person would actually say them aloud in natural spoken Russian/Kazakh. "
|
||||
@@ -174,9 +178,16 @@ def operator_system_prompt(*, language: str, channel_label: str, is_voice: bool,
|
||||
"Paraphrase them naturally instead of quoting them verbatim. "
|
||||
"Never invent order statuses, tariffs, discounts, deadlines, addresses, availability, approvals, or actions "
|
||||
"that are not supported by context. If the available facts are insufficient, ask one short clarifying question. "
|
||||
"After you deliver a complete answer grounded in kb_results (not when you are asking a clarifying question, "
|
||||
"handling an identity/off-topic reply, or closing the call), end reply_text with a brief natural check such as «Ответила ли я на ваш вопрос?» in Russian, or its natural Kazakh equivalent, phrased differently each time so it does not sound scripted. "
|
||||
"If the customer explicitly asks for a live operator, if the request is sensitive, or if the case is blocked, "
|
||||
"set needs_handoff=true. "
|
||||
f"{delivery_hint} "
|
||||
"Return only a JSON object with keys: language, intent, reply_text, extracted_name, confidence, needs_handoff, "
|
||||
"handoff_reason, case_action, kb_refs. case_action must be one of none, close, escalate, keep_open."
|
||||
"handoff_reason, case_action, kb_refs. case_action must be one of none, close, escalate, keep_open. "
|
||||
"For `intent`: if your reply is grounded in one of the provided kb_results, set intent to that snippet's "
|
||||
"intent_code exactly as given (do not translate, reformat, or invent your own code). If no kb_results were "
|
||||
"used, use one of these fixed values as appropriate: identity_question, handoff_request, sensitive_request, "
|
||||
"resolution_confirmed, clarification, kb_answer, unknown. Never invent a new intent value outside of these "
|
||||
"two sources — the platform discards anything else."
|
||||
)
|
||||
|
||||
@@ -603,20 +603,42 @@ def _voice_has_name_correction(text: str | None) -> bool:
|
||||
return any(marker in normalized for marker in correction_markers)
|
||||
|
||||
|
||||
def _voice_reply_with_name(language: str, reply_text: str, name: str | None) -> str:
|
||||
def _voice_reply_already_names_customer(reply_text: str, normalized_short: str) -> bool:
|
||||
tokens = _voice_text_key(reply_text).split(" ")
|
||||
if normalized_short in tokens:
|
||||
return True
|
||||
stem_len = max(len(normalized_short) - 2, 3)
|
||||
stem = normalized_short[:stem_len]
|
||||
return any(len(token) >= stem_len and token.startswith(stem) for token in tokens)
|
||||
|
||||
|
||||
def _voice_greeting_word(language: str) -> str:
|
||||
if str(language or "").strip().lower() == "kz":
|
||||
return "Сәлеметсіз бе"
|
||||
return "Здравствуйте"
|
||||
|
||||
|
||||
def _voice_reply_with_name(
|
||||
language: str,
|
||||
reply_text: str,
|
||||
name: str | None,
|
||||
*,
|
||||
greet: bool = False,
|
||||
) -> str:
|
||||
short_name = _voice_short_name(name)
|
||||
if not short_name:
|
||||
return reply_text
|
||||
prefix = _voice_disclosure_prefix(language)
|
||||
normalized_short = _voice_text_key(short_name)
|
||||
lead = f"{_voice_greeting_word(language)}, {short_name}" if greet else short_name
|
||||
if reply_text.startswith(prefix):
|
||||
rest = reply_text[len(prefix) :].lstrip()
|
||||
if normalized_short in _voice_text_key(rest).split(" "):
|
||||
if _voice_reply_already_names_customer(rest, normalized_short):
|
||||
return reply_text
|
||||
return f"{prefix}{short_name}, {rest}"
|
||||
if normalized_short in _voice_text_key(reply_text).split(" "):
|
||||
return f"{prefix}{lead}, {rest}"
|
||||
if _voice_reply_already_names_customer(reply_text, normalized_short):
|
||||
return reply_text
|
||||
return f"{short_name}, {reply_text}"
|
||||
return f"{lead}, {reply_text}"
|
||||
|
||||
|
||||
def _voice_name_metadata(
|
||||
@@ -746,13 +768,7 @@ def _voice_downstream_name_update(
|
||||
action = "provide" if not name_value else "confirm"
|
||||
candidate = candidate or name_value
|
||||
elif status == "name_obtained":
|
||||
if (
|
||||
candidate
|
||||
and candidate_key
|
||||
and candidate_key != current_key
|
||||
and _voice_has_name_correction(transcript_text)
|
||||
):
|
||||
action = "correct"
|
||||
pass
|
||||
else:
|
||||
if explicit_candidate:
|
||||
action = "provide"
|
||||
@@ -1039,7 +1055,6 @@ def _voice_has_service_topic(text: str | None) -> bool:
|
||||
"ошибк",
|
||||
"проблем",
|
||||
"сбой",
|
||||
"интернет",
|
||||
"связь",
|
||||
"оператор",
|
||||
"менеджер",
|
||||
@@ -1047,6 +1062,40 @@ def _voice_has_service_topic(text: str | None) -> bool:
|
||||
"компан",
|
||||
"подключ",
|
||||
"доставк",
|
||||
"газ",
|
||||
"счётчик",
|
||||
"счетчик",
|
||||
"ваучер",
|
||||
"отключ",
|
||||
"квитанц",
|
||||
"показан",
|
||||
"приложен",
|
||||
"безопасн",
|
||||
"техническ",
|
||||
"плит",
|
||||
"котл",
|
||||
"труб",
|
||||
"утечк",
|
||||
"запах",
|
||||
"абонент",
|
||||
"договор",
|
||||
"поверк",
|
||||
"монтаж",
|
||||
"счёт",
|
||||
"счет",
|
||||
"долг",
|
||||
"задолжен",
|
||||
"перерасчёт",
|
||||
"перерасчет",
|
||||
"регион",
|
||||
"аимак",
|
||||
"aimaq",
|
||||
"qazaqgaz",
|
||||
"казахгаз",
|
||||
"есептегіш",
|
||||
"төлеу",
|
||||
"өтінім",
|
||||
"шарт",
|
||||
)
|
||||
return any(marker in normalized for marker in service_markers)
|
||||
|
||||
@@ -1080,6 +1129,37 @@ def _voice_is_off_domain_request(text: str | None) -> bool:
|
||||
"java",
|
||||
"javascript",
|
||||
"погод",
|
||||
"плов",
|
||||
"приготов",
|
||||
"кулинар",
|
||||
"блюдо",
|
||||
"фильм",
|
||||
"кино",
|
||||
"спорт",
|
||||
"футбол",
|
||||
"хоккей",
|
||||
"теннис",
|
||||
"баскетбол",
|
||||
"президент",
|
||||
"политик",
|
||||
"выборы",
|
||||
"правительств",
|
||||
"курс валют",
|
||||
"криптовалют",
|
||||
"биткоин",
|
||||
"песн",
|
||||
"музык",
|
||||
"танц",
|
||||
"шутк",
|
||||
"загадк",
|
||||
"сериал",
|
||||
"книга",
|
||||
"стихотвор",
|
||||
"чемпионат",
|
||||
"ауа райы",
|
||||
"аспаздық",
|
||||
"кітап",
|
||||
"ән айт",
|
||||
)
|
||||
if any(marker in normalized for marker in broad_markers):
|
||||
return True
|
||||
@@ -1089,7 +1169,14 @@ def _voice_is_off_domain_request(text: str | None) -> bool:
|
||||
"что такое",
|
||||
"объясни",
|
||||
"расскажи про",
|
||||
"почему",
|
||||
"посоветуй",
|
||||
"как приготовить",
|
||||
"какая погода",
|
||||
"кто президент",
|
||||
"кто выиграл",
|
||||
"какая команда",
|
||||
"ауа райы қалай",
|
||||
"әнді айт",
|
||||
)
|
||||
return any(normalized.startswith(prefix) for prefix in broad_openers)
|
||||
|
||||
@@ -1097,13 +1184,12 @@ def _voice_is_off_domain_request(text: str | None) -> bool:
|
||||
def _voice_off_domain_reply(language: str) -> tuple[str, str]:
|
||||
if language == "kz":
|
||||
return (
|
||||
"Men kompaniyamyzdyn qyzmetteri men otinishteri boiynsha komek bere alamyn. "
|
||||
"Eger suraq bizdin qyzmetke qatysty bolsa, qysqasha naqtylaңыз. Qalasaңыз, operatorga qosamyn.",
|
||||
"AI qongyraudyn taqyrybyn kompaniya qyzmetteri sheginde naqtylaudy usyndy.",
|
||||
"Кешіріңіз, мен тек газ қызметтері бойынша көмектесемін: төлеу, есептегіш, ваучер. Қалай көмектесе аламын?",
|
||||
"AI off-topic сұрауды газ қызметтері тақырыбына шектеді.",
|
||||
)
|
||||
return (
|
||||
"Я помогу по вопросам наших услуг и обращений. Если вопрос связан с нашей компанией, скажите коротко, что именно нужно. Если хотите, сразу соединю с оператором.",
|
||||
"AI мягко вернул разговор к вопросам компании и предложил перевод на оператора.",
|
||||
"Извините, я помогаю только по вопросам газа: оплата, счётчики, ваучеры. Чем могу помочь?",
|
||||
"AI отклонил off-topic вопрос и ограничил тему газоснабжением.",
|
||||
)
|
||||
|
||||
|
||||
@@ -1273,7 +1359,13 @@ def _voice_v2_enabled(metadata: dict[str, Any] | None = None) -> bool:
|
||||
return _voice_policy_mode() in {"v2_fast_conversational", "v2_streaming_duplex"}
|
||||
|
||||
|
||||
def _voice_early_intent_bucket(text: str) -> str:
|
||||
def _voice_ack_topic_bucket(text: str) -> str:
|
||||
"""Coarse keyword heuristic used only to pick an ack phrase / early clarifying
|
||||
question (see _voice_ack_kind_for_intent and _voice_early_plan) while the real,
|
||||
KB-grounded decision is still in flight. This is NOT the canonical FAQ intent
|
||||
(see services.shared.intents) and must never be echoed back as the final
|
||||
decision's `intent` value.
|
||||
"""
|
||||
normalized = " ".join(str(text or "").strip().lower().split())
|
||||
if not normalized:
|
||||
return "unknown"
|
||||
@@ -1423,7 +1515,7 @@ def _voice_v2_metadata(
|
||||
if not _voice_v2_enabled(request_metadata):
|
||||
return {}
|
||||
payload = request_metadata if isinstance(request_metadata, dict) else {}
|
||||
early_intent = _voice_early_intent_bucket(transcript_text)
|
||||
early_intent = _voice_ack_topic_bucket(transcript_text)
|
||||
metadata: dict[str, Any] = {
|
||||
"voice_v2_enabled": True,
|
||||
"early_intent": early_intent,
|
||||
@@ -1582,6 +1674,7 @@ def _voice_llm_prompt_messages(
|
||||
"article_id": article.article_id,
|
||||
"title": article.title,
|
||||
"snippet": app._article_snippet(article, limit=240),
|
||||
"intent_code": getattr(article, "intent_code", None),
|
||||
}
|
||||
for article in kb_results[:3]
|
||||
]
|
||||
@@ -1613,6 +1706,12 @@ def _voice_llm_prompt_messages(
|
||||
is_voice=True,
|
||||
config=operator_config,
|
||||
)
|
||||
system_prompt += (
|
||||
" Do not open `reply_text` with a greeting or by addressing the customer by name "
|
||||
"(e.g. do not write 'Здравствуйте, <имя>' or start with '<имя>,'). The system inserts "
|
||||
"the customer's name into the spoken reply separately, so naming them yourself would "
|
||||
"make it get said twice."
|
||||
)
|
||||
if name_status in ("name_not_obtained", "name_followup_required"):
|
||||
system_prompt += (
|
||||
" The user's name is not yet obtained. If the user explicitly provided their name in this turn, "
|
||||
@@ -1675,7 +1774,13 @@ def _voice_llm_decision(
|
||||
)
|
||||
except Exception:
|
||||
return None
|
||||
decision = app._sanitize_decision(raw, fallback_language=language)
|
||||
decision = app._sanitize_decision(
|
||||
raw,
|
||||
fallback_language=language,
|
||||
known_topic_codes=[
|
||||
code for article in kb_results if (code := getattr(article, "intent_code", None))
|
||||
],
|
||||
)
|
||||
if not str(decision.get("reply_text") or "").strip():
|
||||
decision["reply_text"] = _voice_generic_prompt(language)
|
||||
if decision.get("needs_handoff") and not decision.get("handoff_reason"):
|
||||
@@ -1810,7 +1915,7 @@ def _voice_decision(
|
||||
"latency_ms": 1,
|
||||
}
|
||||
|
||||
if not kb_results and _voice_is_off_domain_request(normalized):
|
||||
if _voice_is_off_domain_request(normalized):
|
||||
reply_text, summary_text = _voice_off_domain_reply(language)
|
||||
decision = {
|
||||
"language": language,
|
||||
@@ -2500,10 +2605,12 @@ def turn_voice_session(session_id: str, payload: VoiceAITurnIn) -> VoiceAITurnDe
|
||||
decision_metadata.update(decision.get("metadata") or {})
|
||||
suppress_name_prefix = bool(request_metadata.get("suppress_name_prefix")) if isinstance(request_metadata, dict) else False
|
||||
if effective_name_status == "name_obtained" and effective_name_value and not early_plan_only and not suppress_name_prefix:
|
||||
just_learned_name = current_name_status != "name_obtained"
|
||||
decision["reply_text"] = _voice_reply_with_name(
|
||||
decision["language"],
|
||||
decision["reply_text"],
|
||||
effective_name_value,
|
||||
greet=just_learned_name,
|
||||
)
|
||||
elif inline_name_followup and not decision["needs_handoff"] and not early_plan_only:
|
||||
inline_followup = _voice_inline_name_followup(decision["language"], config)
|
||||
|
||||
@@ -822,8 +822,15 @@ def _media_registration_from_row(row: VoiceAISessionRow, *, queue_code: str | No
|
||||
voice_v2_partial_asr=bool(voice_v2_for_session and _voice_v2_partial_asr_enabled()),
|
||||
voice_v2_duplex=bool(voice_v2_for_session and _voice_v2_duplex_enabled()),
|
||||
voice_v2_streaming_asr_backend=streaming_backend,
|
||||
voice_v2_prebaked_ack=bool(voice_v2_for_session and _voice_v2_prebaked_ack_enabled()),
|
||||
voice_v2_emotive_ack=bool(voice_v2_for_session and _voice_v2_emotive_ack_enabled()),
|
||||
# voice_v2_prebaked_ack and voice_v2_emotive_ack are deliberately NOT gated on
|
||||
# voice_v2_for_session: filler-ack synthesis/caching and phrase-variant rotation
|
||||
# only need a text pool + the ack bank, not the v2 duplex/partial-ASR pipeline.
|
||||
# Without this, calls outside the v2 queue allowlist would synthesize every
|
||||
# filler live via ElevenLabs before it could play — adding real TTS round-trip
|
||||
# time to the one phrase whose whole job is to hide that latency — and would
|
||||
# always get the single fixed "Секунду." fallback string instead of rotating.
|
||||
voice_v2_prebaked_ack=bool(_voice_v2_prebaked_ack_enabled()),
|
||||
voice_v2_emotive_ack=bool(_voice_v2_emotive_ack_enabled()),
|
||||
voice_v2_emotive_ack_ru_only=bool(_voice_v2_emotive_ack_ru_only()),
|
||||
)
|
||||
|
||||
@@ -1328,6 +1335,27 @@ def _process_voice_ai_turn_sync(
|
||||
payload: VoiceAITurnIn,
|
||||
*,
|
||||
auto_handoff: bool,
|
||||
) -> VoiceAITurnDecisionOut:
|
||||
last_exc: Exception | None = None
|
||||
for attempt in range(2):
|
||||
try:
|
||||
return _process_voice_ai_turn_sync_once(session_id, payload, auto_handoff=auto_handoff)
|
||||
except HTTPException as exc:
|
||||
if attempt == 0 and exc.status_code == 502 and "deadlock detected" in str(exc.detail).lower():
|
||||
logging.getLogger(__name__).warning("voice_turn_deadlock_retry session_id=%s", session_id)
|
||||
last_exc = exc
|
||||
continue
|
||||
raise
|
||||
if last_exc is not None:
|
||||
raise last_exc
|
||||
raise RuntimeError("unreachable")
|
||||
|
||||
|
||||
def _process_voice_ai_turn_sync_once(
|
||||
session_id: str,
|
||||
payload: VoiceAITurnIn,
|
||||
*,
|
||||
auto_handoff: bool,
|
||||
) -> VoiceAITurnDecisionOut:
|
||||
session = get_session()
|
||||
voice_session = None
|
||||
|
||||
@@ -482,7 +482,7 @@ class AudioSocketMediaRuntime:
|
||||
if ack_kind == "handoff":
|
||||
return (
|
||||
"Угу, секунду.",
|
||||
"Понял вас, соединяю.",
|
||||
"Понялa вас, соединяю.",
|
||||
"Мхм, соединяю.",
|
||||
"Хорошо, сейчас соединю.",
|
||||
"Да, сейчас соединю.",
|
||||
@@ -492,7 +492,7 @@ class AudioSocketMediaRuntime:
|
||||
return (
|
||||
"Угу, сейчас подскажу.",
|
||||
"Мхм, сориентирую.",
|
||||
"Ага, понял вас.",
|
||||
"Ага, понялa вас.",
|
||||
"Хм, сейчас уточню.",
|
||||
"Хорошо, сейчас подскажу.",
|
||||
"Ясно, сейчас разберусь.",
|
||||
@@ -509,14 +509,12 @@ class AudioSocketMediaRuntime:
|
||||
"Дайте уточню.",
|
||||
)
|
||||
return (
|
||||
"Ага.",
|
||||
"Ясно, минутку.",
|
||||
"Угу.",
|
||||
"Хорошо, сейчас.",
|
||||
"Так, слушаю.",
|
||||
"Момент.",
|
||||
"Понял вас.",
|
||||
"Хм, сейчас гляну.",
|
||||
"Так, слушаю вас.",
|
||||
"Один момент.",
|
||||
"Поняла вас.",
|
||||
"Хмм, сейчас гляну.",
|
||||
)
|
||||
|
||||
def _select_ack_payload(
|
||||
@@ -958,6 +956,11 @@ class AudioSocketMediaRuntime:
|
||||
bool(partial.is_final),
|
||||
transcript_text[:160],
|
||||
)
|
||||
# Fresh speech content is still arriving, so push the
|
||||
# no-speech silence-timeout deadline forward instead of
|
||||
# interrupting a caller who is actively mid-utterance.
|
||||
if actor.state == "listening":
|
||||
actor.listening_since_monotonic = time.monotonic()
|
||||
actor.partial_transcript = transcript_text
|
||||
self._update_stable_partial_transcript(actor, transcript_text, provider_stable=bool(partial.is_stable or partial.is_final))
|
||||
intent = self._detect_early_intent(transcript_text)
|
||||
|
||||
@@ -412,13 +412,18 @@ def ami_loop(stop_event=None) -> None:
|
||||
raise RuntimeError("AMI connection closed")
|
||||
if not frame:
|
||||
continue
|
||||
if frame.get("Event") != "UserEvent":
|
||||
event_name = frame.get("Event")
|
||||
if event_name == "UserEvent":
|
||||
user_event = str(frame.get("UserEvent") or "").strip()
|
||||
if not user_event.startswith(bridge._ami_prefix()):
|
||||
continue
|
||||
bridge._STATE.set_last_event()
|
||||
bridge._record_ami_payload(frame)
|
||||
elif event_name in {"DialEnd", "Hangup"}:
|
||||
bridge._STATE.set_last_event()
|
||||
bridge._record_ami_payload(frame, event_name=event_name)
|
||||
else:
|
||||
continue
|
||||
user_event = str(frame.get("UserEvent") or "").strip()
|
||||
if not user_event.startswith(bridge._ami_prefix()):
|
||||
continue
|
||||
bridge._STATE.set_last_event()
|
||||
bridge._record_ami_payload(frame)
|
||||
except Exception as exc:
|
||||
bridge._STATE.set_connected(False)
|
||||
bridge._STATE.set_error(str(exc))
|
||||
|
||||
@@ -262,6 +262,7 @@ def _start_background_threads() -> None:
|
||||
extra_loops: list[tuple[str, Callable[..., Any]]] = []
|
||||
if _ivr_fastagi_enabled():
|
||||
extra_loops.append(("asterisk-ivr-fastagi-loop", _ivr_fastagi_loop))
|
||||
extra_loops.append(("agent-acw-sweep-loop", _acw_sweep_loop))
|
||||
bridge_runtime.start_background_threads(
|
||||
ami_loop=_ami_loop,
|
||||
failed_retry_loop=_failed_retry_loop,
|
||||
@@ -305,6 +306,9 @@ _update_voice_ai_call_state = bridge_voice_ai.update_call_ai_state
|
||||
_voice_ai_summary_for_call = bridge_voice_ai.voice_ai_summary_for_call
|
||||
_create_escalation = bridge_voice_ai.create_escalation
|
||||
_release_routing_agent = bridge_voice_ai.release_routing_agent
|
||||
_retry_escalation_no_answer = bridge_voice_ai.retry_escalation_no_answer
|
||||
_routing_release_by_agent_id = bridge_voice_ai.routing_release_by_agent_id
|
||||
_set_routing_agent_status = bridge_voice_ai.set_routing_agent_status
|
||||
|
||||
_first_non_empty = bridge_ami.first_non_empty
|
||||
_extract_call_id = bridge_ami.extract_call_id
|
||||
@@ -373,6 +377,8 @@ _process_audio_bridge_ended = bridge_processing.process_audio_bridge_ended
|
||||
_process_call_ended = bridge_processing.process_call_ended
|
||||
_process_operator_connected = bridge_processing.process_operator_connected
|
||||
_process_recording_ready = bridge_processing.process_recording_ready
|
||||
_process_agent_dial_outcome = bridge_processing.process_agent_dial_outcome
|
||||
_acw_sweep_loop = bridge_processing.acw_sweep_loop
|
||||
_process_bridge_row = bridge_processing.process_bridge_row
|
||||
_record_ami_payload = bridge_processing.record_ami_payload
|
||||
_retry_failed_events_once = bridge_processing.retry_failed_events_once
|
||||
|
||||
@@ -4,6 +4,7 @@ from datetime import datetime, timezone
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
from pathlib import Path
|
||||
import threading
|
||||
from typing import Any
|
||||
@@ -884,18 +885,34 @@ def process_call_ended(
|
||||
open_escalation = session.execute(
|
||||
select(EscalationRow).where(
|
||||
EscalationRow.call_id == row.call_id,
|
||||
EscalationRow.status.in_(["requested", "ringing"]),
|
||||
EscalationRow.status.in_(["requested", "ringing", "connected"]),
|
||||
)
|
||||
).scalar_one_or_none()
|
||||
was_talking = open_escalation is not None and open_escalation.status == "connected"
|
||||
if open_escalation is not None:
|
||||
open_escalation.status = "completed" if answered else "failed"
|
||||
open_escalation.completed_at = now
|
||||
if answered and not open_escalation.connected_at:
|
||||
open_escalation.connected_at = now
|
||||
try:
|
||||
bridge._release_routing_agent(row.call_id)
|
||||
except Exception:
|
||||
pass
|
||||
if was_talking and open_escalation and open_escalation.real_agent_id:
|
||||
try:
|
||||
bridge._set_routing_agent_status(open_escalation.real_agent_id, "AFTER_CALL_WORK")
|
||||
except Exception:
|
||||
logger.warning("bridge.agent_acw_status_failed call_id=%s agent_id=%s", row.call_id, open_escalation.real_agent_id)
|
||||
else:
|
||||
try:
|
||||
bridge._release_routing_agent(row.call_id)
|
||||
except Exception:
|
||||
pass
|
||||
if open_escalation is not None:
|
||||
try:
|
||||
bridge._append_interaction_timeline(
|
||||
interaction_id=link.interaction_id,
|
||||
action="escalation.completed" if answered else "escalation.failed_client_disconnected",
|
||||
metadata={"call_id": row.call_id, "escalation_id": open_escalation.escalation_id, "agent_id": open_escalation.real_agent_id},
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
try:
|
||||
bridge._notify_voice_ai_telephony_event(
|
||||
voice_session_id=link.voice_session_id,
|
||||
@@ -969,6 +986,48 @@ def process_operator_connected(
|
||||
link.telephony_status = "connected"
|
||||
link.connected_at = now
|
||||
link.updated_at = now
|
||||
|
||||
open_escalation = session.execute(
|
||||
select(EscalationRow).where(
|
||||
EscalationRow.call_id == row.call_id,
|
||||
EscalationRow.status == "ringing",
|
||||
)
|
||||
).scalar_one_or_none()
|
||||
if open_escalation is not None:
|
||||
open_escalation.status = "connected"
|
||||
open_escalation.connected_at = now
|
||||
if open_escalation.real_agent_id:
|
||||
try:
|
||||
bridge._set_routing_agent_status(open_escalation.real_agent_id, "TALKING")
|
||||
except Exception:
|
||||
logger.warning("bridge.agent_talking_status_failed call_id=%s agent_id=%s", row.call_id, open_escalation.real_agent_id)
|
||||
try:
|
||||
bridge._append_interaction_timeline(
|
||||
interaction_id=link.interaction_id,
|
||||
action="escalation.agent_connected",
|
||||
metadata={"call_id": row.call_id, "agent_id": open_escalation.real_agent_id, "escalation_id": open_escalation.escalation_id},
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
try:
|
||||
bridge._emit_voice_event(
|
||||
event_type="AgentConnected",
|
||||
call_id=row.call_id,
|
||||
interaction_id=link.interaction_id,
|
||||
payload={"escalation_id": open_escalation.escalation_id, "agent_id": open_escalation.real_agent_id},
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
try:
|
||||
bridge._emit_voice_event(
|
||||
event_type="TransferCompleted",
|
||||
call_id=row.call_id,
|
||||
interaction_id=link.interaction_id,
|
||||
payload={"escalation_id": open_escalation.escalation_id, "agent_id": open_escalation.real_agent_id},
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
_upsert_voice_reporting_fact(
|
||||
session,
|
||||
link,
|
||||
@@ -1173,6 +1232,55 @@ def process_recording_ready(
|
||||
local_path.unlink(missing_ok=True)
|
||||
|
||||
|
||||
_NO_ANSWER_DIAL_STATUSES = {"NOANSWER", "BUSY", "CANCEL", "CHANUNAVAIL", "CONGESTION"}
|
||||
_NO_ANSWER_HANGUP_CAUSES = {"1", "3", "17", "18", "19", "20", "21", "22", "34", "38"}
|
||||
|
||||
|
||||
def process_agent_dial_outcome(session, row: AsteriskEventLogRow, payload: dict[str, Any]) -> None:
|
||||
bridge = _bridge_app()
|
||||
call_id = row.call_id
|
||||
if not call_id:
|
||||
bridge._mark_log_forwarded(session, row, interaction_id=None)
|
||||
return
|
||||
|
||||
dial_status = str(payload.get("DialStatus") or "").strip().upper()
|
||||
hangup_cause = str(payload.get("Cause") or "").strip()
|
||||
|
||||
is_no_answer_outcome = (
|
||||
(row.ami_event_name == "DialEnd" and dial_status in _NO_ANSWER_DIAL_STATUSES)
|
||||
or (row.ami_event_name == "Hangup" and hangup_cause in _NO_ANSWER_HANGUP_CAUSES)
|
||||
)
|
||||
if not is_no_answer_outcome:
|
||||
bridge._mark_log_forwarded(session, row, interaction_id=None)
|
||||
return
|
||||
|
||||
escalation = session.execute(
|
||||
select(EscalationRow).where(
|
||||
EscalationRow.call_id == call_id,
|
||||
EscalationRow.status == "ringing",
|
||||
)
|
||||
).scalar_one_or_none()
|
||||
if escalation is None:
|
||||
bridge._mark_log_forwarded(session, row, interaction_id=None)
|
||||
return
|
||||
|
||||
dialed_channel = payload.get("DestChannel") if row.ami_event_name == "DialEnd" else payload.get("Channel")
|
||||
dialed_extension = bridge._extract_extension_from_channel(dialed_channel)
|
||||
if dialed_extension and escalation.assigned_agent_id and dialed_extension != escalation.assigned_agent_id:
|
||||
bridge._mark_log_forwarded(session, row, interaction_id=None)
|
||||
return
|
||||
|
||||
try:
|
||||
bridge._retry_escalation_no_answer(
|
||||
session,
|
||||
call_id=call_id,
|
||||
dial_outcome=dial_status or f"hangup_cause_{hangup_cause}",
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("bridge.agent_dial_outcome_retry_failed call_id=%s", call_id)
|
||||
bridge._mark_log_forwarded(session, row, interaction_id=None)
|
||||
|
||||
|
||||
def process_bridge_row(session, row: AsteriskEventLogRow) -> AsteriskEventLogRow:
|
||||
bridge = _bridge_app()
|
||||
payload = json.loads(row.payload_json or "{}")
|
||||
@@ -1209,6 +1317,8 @@ def process_bridge_row(session, row: AsteriskEventLogRow) -> AsteriskEventLogRow
|
||||
bridge._process_call_ended(session, row, payload)
|
||||
elif row.ami_event_name == f"{bridge._ami_prefix()}RecordingReady":
|
||||
bridge._process_recording_ready(session, row, payload)
|
||||
elif row.ami_event_name in {"DialEnd", "Hangup"}:
|
||||
bridge._process_agent_dial_outcome(session, row, payload)
|
||||
else:
|
||||
row.forward_status = "received"
|
||||
row.updated_at = utc_now_iso()
|
||||
@@ -1216,9 +1326,9 @@ def process_bridge_row(session, row: AsteriskEventLogRow) -> AsteriskEventLogRow
|
||||
return row
|
||||
|
||||
|
||||
def record_ami_payload(payload: dict[str, Any]) -> None:
|
||||
def record_ami_payload(payload: dict[str, Any], *, event_name: str | None = None) -> None:
|
||||
bridge = _bridge_app()
|
||||
event_name = str(payload.get("UserEvent") or "").strip()
|
||||
event_name = str(event_name or payload.get("UserEvent") or "").strip()
|
||||
call_id = bridge._extract_call_id(payload)
|
||||
linked_id = bridge._extract_linked_id(payload, call_id)
|
||||
if not event_name or not call_id:
|
||||
@@ -1262,6 +1372,42 @@ def retry_failed_events_once() -> None:
|
||||
bridge._process_claimed_bridge_event(bridge_event_id)
|
||||
|
||||
|
||||
def acw_duration_seconds() -> int:
|
||||
raw = str(os.environ.get("ACW_DURATION_SECONDS", "30")).strip()
|
||||
try:
|
||||
return max(int(raw), 1)
|
||||
except ValueError:
|
||||
return 30
|
||||
|
||||
|
||||
def acw_sweep_interval_seconds() -> float:
|
||||
return max(min(float(acw_duration_seconds()) / 2, 15.0), 5.0)
|
||||
|
||||
|
||||
def acw_sweep_once() -> None:
|
||||
from services.routing_service import engine as routing_engine
|
||||
|
||||
session = get_session()
|
||||
try:
|
||||
swept = routing_engine.sweep_after_call_work(session, older_than_seconds=acw_duration_seconds())
|
||||
for agent in swept:
|
||||
logger.warning("bridge.agent_acw_swept agent_id=%s", agent.agent_id)
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
|
||||
def acw_sweep_loop(stop_event: threading.Event | None = None) -> None:
|
||||
bridge = _bridge_app()
|
||||
active_stop_event = stop_event or bridge._background_stop_event()
|
||||
while not active_stop_event.is_set():
|
||||
try:
|
||||
acw_sweep_once()
|
||||
except Exception:
|
||||
logger.exception("bridge.acw_sweep_failed")
|
||||
if active_stop_event.wait(acw_sweep_interval_seconds()):
|
||||
break
|
||||
|
||||
|
||||
def failed_retry_loop(stop_event: threading.Event | None = None) -> None:
|
||||
bridge = _bridge_app()
|
||||
active_stop_event = stop_event or bridge._background_stop_event()
|
||||
|
||||
@@ -63,6 +63,57 @@ def append_interaction_timeline(
|
||||
)
|
||||
|
||||
|
||||
def _append_escalation_timeline(
|
||||
escalation: EscalationRow,
|
||||
link: AsteriskCallLinkRow,
|
||||
*,
|
||||
action: str,
|
||||
extra: dict[str, Any] | None = None,
|
||||
) -> None:
|
||||
try:
|
||||
append_interaction_timeline(
|
||||
interaction_id=link.interaction_id,
|
||||
action=action,
|
||||
metadata={
|
||||
"call_id": escalation.call_id,
|
||||
"escalation_id": escalation.escalation_id,
|
||||
"from_level": escalation.from_level,
|
||||
"to_level": escalation.to_level,
|
||||
"attempt_count": escalation.attempt_count,
|
||||
**(extra or {}),
|
||||
},
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def _emit_escalation_event(
|
||||
*,
|
||||
event_type: str,
|
||||
escalation: EscalationRow,
|
||||
link: AsteriskCallLinkRow,
|
||||
extra: dict[str, Any] | None = None,
|
||||
) -> None:
|
||||
bridge = _bridge_app()
|
||||
try:
|
||||
bridge._emit_voice_event(
|
||||
event_type=event_type,
|
||||
call_id=escalation.call_id,
|
||||
interaction_id=link.interaction_id,
|
||||
payload={
|
||||
"escalation_id": escalation.escalation_id,
|
||||
"tenant_id": escalation.tenant_id,
|
||||
"from_level": escalation.from_level,
|
||||
"to_level": escalation.to_level,
|
||||
"reason_code": escalation.reason_code,
|
||||
"attempt_count": escalation.attempt_count,
|
||||
**(extra or {}),
|
||||
},
|
||||
)
|
||||
except Exception:
|
||||
LOGGER.warning("bridge.escalation_event_failed event_type=%s call_id=%s", event_type, escalation.call_id)
|
||||
|
||||
|
||||
def start_voice_ai_session(
|
||||
*,
|
||||
call_id: str,
|
||||
@@ -188,7 +239,8 @@ def _reserve_routing_agent(
|
||||
level: str,
|
||||
tenant_id: str | None,
|
||||
required_skills: list[str] | None = None,
|
||||
) -> str | None:
|
||||
exclude_agent_ids: list[str] | None = None,
|
||||
) -> dict[str, Any] | None:
|
||||
bridge = _bridge_app()
|
||||
try:
|
||||
response = bridge._post_json(
|
||||
@@ -198,7 +250,7 @@ def _reserve_routing_agent(
|
||||
"level": level,
|
||||
"tenant_id": tenant_id,
|
||||
"required_skills": required_skills or [],
|
||||
"exclude_agent_ids": [],
|
||||
"exclude_agent_ids": exclude_agent_ids or [],
|
||||
},
|
||||
timeout_seconds=bridge._callcontrol_side_effect_timeout_seconds(),
|
||||
max_attempts=1,
|
||||
@@ -212,8 +264,43 @@ def _reserve_routing_agent(
|
||||
tenant_id,
|
||||
)
|
||||
return None
|
||||
agent_id = str((response or {}).get("agent_id") or "").strip()
|
||||
extension = str((response or {}).get("extension") or "").strip()
|
||||
return extension or None
|
||||
if not agent_id or not extension:
|
||||
return None
|
||||
return {
|
||||
"agent_id": agent_id,
|
||||
"extension": extension,
|
||||
"endpoint": (response or {}).get("endpoint"),
|
||||
"display_name": (response or {}).get("display_name"),
|
||||
}
|
||||
|
||||
|
||||
def set_routing_agent_status(agent_id: str, status: str) -> None:
|
||||
bridge = _bridge_app()
|
||||
try:
|
||||
bridge._patch_json(
|
||||
f"{bridge._routing_service_url()}/internal/routing/agents/{agent_id}/status",
|
||||
{"status": status},
|
||||
timeout_seconds=bridge._callcontrol_side_effect_timeout_seconds(),
|
||||
max_attempts=1,
|
||||
retry_backoff_seconds=0.0,
|
||||
)
|
||||
except Exception:
|
||||
LOGGER.warning("bridge.routing_set_status_failed agent_id=%s status=%s", agent_id, status)
|
||||
|
||||
|
||||
def _redirect_channel_to_agent(*, channel: str, extension: str) -> None:
|
||||
bridge = _bridge_app()
|
||||
bridge._ami_action(
|
||||
"Redirect",
|
||||
{
|
||||
"Channel": channel,
|
||||
"Context": bridge._transfer_context(),
|
||||
"Exten": extension,
|
||||
"Priority": 1,
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
def _resolve_handoff_extension(
|
||||
@@ -224,7 +311,7 @@ def _resolve_handoff_extension(
|
||||
target_level: str | None = None,
|
||||
tenant_id: str | None = None,
|
||||
required_skills: list[str] | None = None,
|
||||
) -> tuple[str, str, str | None, str | None]:
|
||||
) -> tuple[str, str, str | None, str | None, str | None]:
|
||||
bridge = _bridge_app()
|
||||
queue_code = _queue_code_for_queue_id(target_queue_id) or str(fallback_queue_code or "").strip()
|
||||
if not queue_code:
|
||||
@@ -233,14 +320,14 @@ def _resolve_handoff_extension(
|
||||
level = target_level or bridge._routing_level_for_queue_code(queue_code)
|
||||
if level and call_id:
|
||||
resolved_tenant_id = tenant_id if tenant_id is not None else bridge._routing_tenant_for_queue_code(queue_code)
|
||||
reserved_extension = _reserve_routing_agent(
|
||||
reserved_agent = _reserve_routing_agent(
|
||||
call_id=call_id,
|
||||
level=level,
|
||||
tenant_id=resolved_tenant_id,
|
||||
required_skills=required_skills,
|
||||
)
|
||||
if reserved_extension:
|
||||
return queue_code, reserved_extension, level, resolved_tenant_id
|
||||
if reserved_agent:
|
||||
return queue_code, reserved_agent["extension"], level, resolved_tenant_id, reserved_agent["agent_id"]
|
||||
raise HTTPException(status_code=409, detail=f"No available {level} agent right now")
|
||||
|
||||
extension = bridge._transfer_target_map().get(queue_code)
|
||||
@@ -249,7 +336,7 @@ def _resolve_handoff_extension(
|
||||
status_code=400,
|
||||
detail=f"Unknown transfer queue_code: {queue_code}",
|
||||
)
|
||||
return queue_code, extension, None, None
|
||||
return queue_code, extension, None, None, None
|
||||
|
||||
|
||||
def _resolve_handoff_channel(session, link: AsteriskCallLinkRow) -> str:
|
||||
@@ -459,11 +546,33 @@ def request_handoff(call_id: str, body: VoiceAIHandoffRequestIn, actor: dict) ->
|
||||
if voice_session is None:
|
||||
raise HTTPException(status_code=404, detail="Voice AI session not found")
|
||||
|
||||
queue_code, target_extension, resolved_level, resolved_tenant_id = _resolve_handoff_extension(
|
||||
queue_code, target_extension, resolved_level, resolved_tenant_id, reserved_agent_id = _resolve_handoff_extension(
|
||||
body.target_queue_id or voice_session.handoff_target_queue_id or link.queue_id,
|
||||
fallback_queue_code=_queue_code_for_queue_id(link.queue_id),
|
||||
call_id=call_id,
|
||||
)
|
||||
escalation: EscalationRow | None = None
|
||||
if reserved_agent_id and target_extension != "7100":
|
||||
now_reserved = utc_now_iso()
|
||||
escalation = EscalationRow(
|
||||
escalation_id=new_id("esc"),
|
||||
call_id=call_id,
|
||||
tenant_id=resolved_tenant_id,
|
||||
from_level=str(link.current_level or "L1"),
|
||||
to_level=resolved_level,
|
||||
reason_code="AI_HANDOFF",
|
||||
required_skills_json="[]",
|
||||
priority=3,
|
||||
status="ringing",
|
||||
real_agent_id=reserved_agent_id,
|
||||
assigned_agent_id=target_extension,
|
||||
attempted_agent_ids_json=json.dumps([reserved_agent_id], ensure_ascii=False),
|
||||
requested_at=now_reserved,
|
||||
)
|
||||
session.add(escalation)
|
||||
session.flush()
|
||||
_append_escalation_timeline(escalation, link, action="escalation.agent_reserved", extra={"agent_id": reserved_agent_id})
|
||||
_emit_escalation_event(event_type="AgentReserved", escalation=escalation, link=link, extra={"agent_id": reserved_agent_id})
|
||||
handoff_metadata = body.metadata or {}
|
||||
actor_user = str(actor.get("user") or actor.get("sub") or "ai-voice-runtime").strip()
|
||||
actor_role = str(actor.get("role") or "admin").strip() or "admin"
|
||||
@@ -510,15 +619,33 @@ def request_handoff(call_id: str, body: VoiceAIHandoffRequestIn, actor: dict) ->
|
||||
target_extension,
|
||||
body.target_queue_id or voice_session.handoff_target_queue_id or link.queue_id,
|
||||
)
|
||||
ami_result = bridge._ami_action(
|
||||
"Redirect",
|
||||
{
|
||||
"Channel": channel,
|
||||
"Context": bridge._transfer_context(),
|
||||
"Exten": target_extension,
|
||||
"Priority": 1,
|
||||
},
|
||||
)
|
||||
try:
|
||||
ami_result = bridge._ami_action(
|
||||
"Redirect",
|
||||
{
|
||||
"Channel": channel,
|
||||
"Context": bridge._transfer_context(),
|
||||
"Exten": target_extension,
|
||||
"Priority": 1,
|
||||
},
|
||||
)
|
||||
except Exception as exc:
|
||||
if escalation is not None:
|
||||
try:
|
||||
routing_release_by_agent_id(reserved_agent_id)
|
||||
except Exception:
|
||||
LOGGER.warning("bridge.handoff_release_agent_failed call_id=%s agent_id=%s", call_id, reserved_agent_id)
|
||||
escalation.status = "failed"
|
||||
escalation.completed_at = utc_now_iso()
|
||||
session.commit()
|
||||
_append_escalation_timeline(escalation, link, action="escalation.transfer_failed", extra={"agent_id": reserved_agent_id, "error": str(exc)})
|
||||
_emit_escalation_event(event_type="TransferFailed", escalation=escalation, link=link, extra={"agent_id": reserved_agent_id, "error": str(exc)})
|
||||
raise HTTPException(status_code=502, detail=f"Failed to redirect call to agent: {exc}") from exc
|
||||
|
||||
if escalation is not None:
|
||||
set_routing_agent_status(reserved_agent_id, "RINGING")
|
||||
_append_escalation_timeline(escalation, link, action="escalation.agent_ringing", extra={"agent_id": reserved_agent_id})
|
||||
_emit_escalation_event(event_type="AgentRinging", escalation=escalation, link=link, extra={"agent_id": reserved_agent_id})
|
||||
|
||||
now = utc_now_iso()
|
||||
link.voice_session_id = (
|
||||
@@ -666,6 +793,7 @@ def _escalation_to_out(row: EscalationRow) -> EscalationOut:
|
||||
summary=row.summary,
|
||||
status=row.status,
|
||||
assigned_agent_id=row.assigned_agent_id,
|
||||
attempt_count=row.attempt_count or 0,
|
||||
requested_at=row.requested_at,
|
||||
connected_at=row.connected_at,
|
||||
completed_at=row.completed_at,
|
||||
@@ -709,28 +837,32 @@ def create_escalation(call_id: str, body: EscalationRequestIn, actor: dict) -> E
|
||||
session.flush()
|
||||
|
||||
channel = _resolve_handoff_channel(session, link)
|
||||
reserved_extension = _reserve_routing_agent(
|
||||
agent = _reserve_routing_agent(
|
||||
call_id=call_id,
|
||||
level=body.target_level,
|
||||
tenant_id=link.tenant_id,
|
||||
required_skills=body.required_skills,
|
||||
)
|
||||
if not reserved_extension:
|
||||
if not agent:
|
||||
escalation.status = "failed"
|
||||
escalation.completed_at = now
|
||||
session.commit()
|
||||
_emit_escalation_event(
|
||||
event_type="TransferFailed",
|
||||
escalation=escalation,
|
||||
link=link,
|
||||
extra={"reason": "no_available_agent"},
|
||||
)
|
||||
raise HTTPException(status_code=409, detail=f"No available {body.target_level} agent right now")
|
||||
|
||||
escalation.real_agent_id = agent["agent_id"]
|
||||
escalation.attempted_agent_ids_json = json.dumps([agent["agent_id"]], ensure_ascii=False)
|
||||
session.flush()
|
||||
_append_escalation_timeline(escalation, link, action="escalation.agent_reserved", extra={"agent_id": agent["agent_id"]})
|
||||
_emit_escalation_event(event_type="AgentReserved", escalation=escalation, link=link, extra={"agent_id": agent["agent_id"]})
|
||||
|
||||
try:
|
||||
bridge._ami_action(
|
||||
"Redirect",
|
||||
{
|
||||
"Channel": channel,
|
||||
"Context": bridge._transfer_context(),
|
||||
"Exten": reserved_extension,
|
||||
"Priority": 1,
|
||||
},
|
||||
)
|
||||
_redirect_channel_to_agent(channel=channel, extension=agent["extension"])
|
||||
except Exception as exc:
|
||||
try:
|
||||
bridge._release_routing_agent(call_id)
|
||||
@@ -739,10 +871,13 @@ def create_escalation(call_id: str, body: EscalationRequestIn, actor: dict) -> E
|
||||
escalation.status = "failed"
|
||||
escalation.completed_at = utc_now_iso()
|
||||
session.commit()
|
||||
_append_escalation_timeline(escalation, link, action="escalation.transfer_failed", extra={"agent_id": agent["agent_id"], "error": str(exc)})
|
||||
_emit_escalation_event(event_type="TransferFailed", escalation=escalation, link=link, extra={"agent_id": agent["agent_id"], "error": str(exc)})
|
||||
raise HTTPException(status_code=502, detail=f"Failed to redirect call to agent: {exc}") from exc
|
||||
|
||||
set_routing_agent_status(agent["agent_id"], "RINGING")
|
||||
escalation.status = "ringing"
|
||||
escalation.assigned_agent_id = reserved_extension
|
||||
escalation.assigned_agent_id = agent["extension"]
|
||||
link.current_level = body.target_level
|
||||
link.required_skills_json = json.dumps(body.required_skills, ensure_ascii=False)
|
||||
link.priority = body.priority
|
||||
@@ -751,11 +886,118 @@ def create_escalation(call_id: str, body: EscalationRequestIn, actor: dict) -> E
|
||||
link.operator_extension = None
|
||||
link.updated_at = now
|
||||
session.commit()
|
||||
_append_escalation_timeline(escalation, link, action="escalation.agent_ringing", extra={"agent_id": agent["agent_id"]})
|
||||
_emit_escalation_event(event_type="AgentRinging", escalation=escalation, link=link, extra={"agent_id": agent["agent_id"]})
|
||||
return _escalation_to_out(escalation)
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
|
||||
def retry_escalation_no_answer(session, *, call_id: str, dial_outcome: str) -> None:
|
||||
"""AC-08 / ТЗ §14: агент не ответил — освободить его и попробовать следующего.
|
||||
|
||||
Вызывается из bridge_processing при нативном AMI DialEnd/Hangup с исходом
|
||||
NOANSWER/BUSY/CANCEL/CHANUNAVAIL/CONGESTION на канале агента, зарезервированного
|
||||
под активную (status='ringing') эскалацию этого call_id.
|
||||
"""
|
||||
bridge = _bridge_app()
|
||||
escalation = session.execute(
|
||||
select(EscalationRow).where(
|
||||
EscalationRow.call_id == call_id,
|
||||
EscalationRow.status == "ringing",
|
||||
)
|
||||
).scalar_one_or_none()
|
||||
if escalation is None:
|
||||
return
|
||||
|
||||
link = bridge._find_call_link(session, call_id)
|
||||
if link is None:
|
||||
return
|
||||
|
||||
now = utc_now_iso()
|
||||
no_answer_agent_id = escalation.real_agent_id
|
||||
if no_answer_agent_id:
|
||||
try:
|
||||
routing_release_by_agent_id(no_answer_agent_id)
|
||||
except Exception:
|
||||
LOGGER.warning("bridge.escalation_retry_release_failed call_id=%s agent_id=%s", call_id, no_answer_agent_id)
|
||||
|
||||
escalation.attempt_count = (escalation.attempt_count or 0) + 1
|
||||
session.flush()
|
||||
_append_escalation_timeline(
|
||||
escalation,
|
||||
link,
|
||||
action="escalation.agent_no_answer",
|
||||
extra={"agent_id": no_answer_agent_id, "dial_outcome": dial_outcome},
|
||||
)
|
||||
_emit_escalation_event(
|
||||
event_type="AgentNoAnswer",
|
||||
escalation=escalation,
|
||||
link=link,
|
||||
extra={"agent_id": no_answer_agent_id, "dial_outcome": dial_outcome},
|
||||
)
|
||||
|
||||
attempted_ids = json.loads(escalation.attempted_agent_ids_json or "[]")
|
||||
# Use the already-known, actively-maintained channel name directly instead of
|
||||
# _resolve_handoff_channel()'s live AMI CoreShowChannels re-discovery: that
|
||||
# round-trip can take ~10s, which races (and loses) against the dialplan's
|
||||
# own short MusicOnHold-then-hangup wait window for this exact retry path.
|
||||
channel = str(link.channel_name or "").strip() or _resolve_handoff_channel(session, link)
|
||||
required_skills = json.loads(escalation.required_skills_json or "[]")
|
||||
next_agent = _reserve_routing_agent(
|
||||
call_id=call_id,
|
||||
level=escalation.to_level,
|
||||
tenant_id=escalation.tenant_id,
|
||||
required_skills=required_skills,
|
||||
exclude_agent_ids=attempted_ids,
|
||||
)
|
||||
if not next_agent:
|
||||
escalation.status = "failed"
|
||||
escalation.completed_at = now
|
||||
session.commit()
|
||||
_append_escalation_timeline(escalation, link, action="escalation.failed_no_agents")
|
||||
_emit_escalation_event(event_type="TransferFailed", escalation=escalation, link=link, extra={"reason": "no_more_agents"})
|
||||
LOGGER.warning("bridge.escalation_retry_exhausted call_id=%s attempts=%s", call_id, escalation.attempt_count)
|
||||
return
|
||||
|
||||
escalation.real_agent_id = next_agent["agent_id"]
|
||||
escalation.attempted_agent_ids_json = json.dumps(attempted_ids + [next_agent["agent_id"]], ensure_ascii=False)
|
||||
session.flush()
|
||||
_append_escalation_timeline(escalation, link, action="escalation.agent_reserved", extra={"agent_id": next_agent["agent_id"]})
|
||||
_emit_escalation_event(event_type="AgentReserved", escalation=escalation, link=link, extra={"agent_id": next_agent["agent_id"]})
|
||||
|
||||
try:
|
||||
_redirect_channel_to_agent(channel=channel, extension=next_agent["extension"])
|
||||
except Exception as exc:
|
||||
try:
|
||||
routing_release_by_agent_id(next_agent["agent_id"])
|
||||
except Exception:
|
||||
pass
|
||||
escalation.status = "failed"
|
||||
escalation.completed_at = utc_now_iso()
|
||||
session.commit()
|
||||
_append_escalation_timeline(escalation, link, action="escalation.transfer_failed", extra={"agent_id": next_agent["agent_id"], "error": str(exc)})
|
||||
_emit_escalation_event(event_type="TransferFailed", escalation=escalation, link=link, extra={"agent_id": next_agent["agent_id"], "error": str(exc)})
|
||||
return
|
||||
|
||||
set_routing_agent_status(next_agent["agent_id"], "RINGING")
|
||||
escalation.assigned_agent_id = next_agent["extension"]
|
||||
session.commit()
|
||||
_append_escalation_timeline(escalation, link, action="escalation.agent_ringing", extra={"agent_id": next_agent["agent_id"]})
|
||||
_emit_escalation_event(event_type="AgentRinging", escalation=escalation, link=link, extra={"agent_id": next_agent["agent_id"]})
|
||||
|
||||
|
||||
def routing_release_by_agent_id(agent_id: str, *, next_status: str = "AVAILABLE") -> None:
|
||||
bridge = _bridge_app()
|
||||
bridge._post_json(
|
||||
f"{bridge._routing_service_url()}/internal/routing/release-agent",
|
||||
{"agent_id": agent_id, "next_status": next_status},
|
||||
timeout_seconds=bridge._callcontrol_side_effect_timeout_seconds(),
|
||||
max_attempts=1,
|
||||
retry_backoff_seconds=0.0,
|
||||
)
|
||||
|
||||
|
||||
def update_call_ai_state(call_id: str, body: VoiceAICallStateUpdateIn, actor: dict) -> VoiceLiveCallOut:
|
||||
bridge = _bridge_app()
|
||||
assert_trusted_voice_runtime_actor(actor)
|
||||
|
||||
@@ -40,6 +40,7 @@ def _article_out(row: KBArticleRow) -> KBArticleOut:
|
||||
article_id=row.article_id,
|
||||
category_id=row.category_id,
|
||||
article_group_id=resolve_article_group_id(row.article_id, row.article_group_id),
|
||||
intent_code=row.intent_code,
|
||||
language=normalize_kb_language(row.language),
|
||||
title=row.title,
|
||||
body=row.body,
|
||||
@@ -49,6 +50,10 @@ def _article_out(row: KBArticleRow) -> KBArticleOut:
|
||||
)
|
||||
|
||||
|
||||
def _normalize_intent_code(value: str | None) -> str | None:
|
||||
return str(value or "").strip().upper() or None
|
||||
|
||||
|
||||
def _article_group_expr():
|
||||
return func.coalesce(KBArticleRow.article_group_id, KBArticleRow.article_id)
|
||||
|
||||
@@ -133,6 +138,7 @@ def create_article(
|
||||
article_id=article_id,
|
||||
category_id=payload.category_id,
|
||||
article_group_id=article_group_id,
|
||||
intent_code=_normalize_intent_code(payload.intent_code),
|
||||
language=language,
|
||||
title=payload.title,
|
||||
body=payload.body,
|
||||
@@ -196,6 +202,8 @@ def update_article(
|
||||
|
||||
if "article_group_id" in data:
|
||||
row.article_group_id = target_group_id
|
||||
if "intent_code" in data:
|
||||
row.intent_code = _normalize_intent_code(data["intent_code"])
|
||||
if "language" in data:
|
||||
row.language = target_language
|
||||
if "title" in data:
|
||||
|
||||
@@ -18,6 +18,7 @@ from services.shared.models import (
|
||||
QueueOut,
|
||||
RoutingAgentReserveIn,
|
||||
RoutingAgentReserveOut,
|
||||
RoutingAgentStatusIn,
|
||||
)
|
||||
from services.shared.security import require_roles
|
||||
from services.shared.sql_init import init_sql_schema
|
||||
@@ -343,6 +344,7 @@ def _escalation_to_out(row: EscalationRow) -> EscalationOut:
|
||||
summary=row.summary,
|
||||
status=row.status,
|
||||
assigned_agent_id=row.assigned_agent_id,
|
||||
attempt_count=row.attempt_count or 0,
|
||||
requested_at=row.requested_at,
|
||||
connected_at=row.connected_at,
|
||||
completed_at=row.completed_at,
|
||||
@@ -399,11 +401,32 @@ def release_agent_endpoint(
|
||||
_: dict = Depends(require_roles(Role.ADMIN)),
|
||||
) -> dict:
|
||||
call_id = str(payload.get("call_id") or "").strip()
|
||||
if not call_id:
|
||||
raise HTTPException(status_code=400, detail="call_id is required")
|
||||
agent_id = str(payload.get("agent_id") or "").strip()
|
||||
next_status = str(payload.get("next_status") or "AVAILABLE").strip() or "AVAILABLE"
|
||||
if not call_id and not agent_id:
|
||||
raise HTTPException(status_code=400, detail="call_id or agent_id is required")
|
||||
session = get_session()
|
||||
try:
|
||||
agent = routing_engine.release_agent_by_call_id(session, call_id=call_id)
|
||||
return {"call_id": call_id, "released_agent_id": agent.agent_id if agent else None}
|
||||
if agent_id:
|
||||
agent = routing_engine.release_agent_by_id(session, agent_id=agent_id, next_status=next_status)
|
||||
else:
|
||||
agent = routing_engine.release_agent_by_call_id(session, call_id=call_id)
|
||||
return {"call_id": call_id or None, "released_agent_id": agent.agent_id if agent else None}
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
|
||||
@app.patch("/internal/routing/agents/{agent_id}/status")
|
||||
def set_agent_status_endpoint(
|
||||
agent_id: str,
|
||||
payload: RoutingAgentStatusIn,
|
||||
_: dict = Depends(require_roles(Role.ADMIN)),
|
||||
) -> dict:
|
||||
session = get_session()
|
||||
try:
|
||||
agent = routing_engine.set_agent_status(session, agent_id=agent_id, status=payload.status)
|
||||
if agent is None:
|
||||
raise HTTPException(status_code=404, detail="Agent not found")
|
||||
return {"agent_id": agent.agent_id, "status": agent.status}
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import select, text
|
||||
|
||||
@@ -117,6 +118,41 @@ def release_agent_by_call_id(session, *, call_id: str) -> AgentRow | None:
|
||||
return agent
|
||||
|
||||
|
||||
def set_agent_status(session, *, agent_id: str, status: str) -> AgentRow | None:
|
||||
agent = session.execute(
|
||||
select(AgentRow).where(AgentRow.agent_id == agent_id)
|
||||
).scalar_one_or_none()
|
||||
if agent is None:
|
||||
return None
|
||||
agent.status = status
|
||||
agent.updated_at = utc_now_iso()
|
||||
session.commit()
|
||||
return agent
|
||||
|
||||
|
||||
def sweep_after_call_work(session, *, older_than_seconds: int) -> list[AgentRow]:
|
||||
cutoff = utc_now_iso()
|
||||
rows = session.execute(
|
||||
select(AgentRow).where(AgentRow.status == "AFTER_CALL_WORK")
|
||||
).scalars().all()
|
||||
swept: list[AgentRow] = []
|
||||
for agent in rows:
|
||||
try:
|
||||
age_seconds = (
|
||||
datetime.fromisoformat(cutoff) - datetime.fromisoformat(str(agent.updated_at))
|
||||
).total_seconds()
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
if age_seconds >= older_than_seconds:
|
||||
agent.status = "AVAILABLE"
|
||||
agent.current_call_id = None
|
||||
agent.updated_at = cutoff
|
||||
swept.append(agent)
|
||||
if swept:
|
||||
session.commit()
|
||||
return swept
|
||||
|
||||
|
||||
def release_agent_by_id(session, *, agent_id: str, next_status: str = "AVAILABLE") -> AgentRow | None:
|
||||
agent = session.execute(
|
||||
select(AgentRow).where(AgentRow.agent_id == agent_id)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
from typing import Any
|
||||
|
||||
from sqlalchemy import select
|
||||
@@ -13,22 +14,48 @@ from services.shared.sql_models import AIOperatorSettingsRow
|
||||
AI_OPERATOR_SETTINGS_KEY = "global"
|
||||
|
||||
|
||||
def _env_field(name: str, default: str) -> str:
|
||||
value = os.getenv(name)
|
||||
if value is None:
|
||||
return default
|
||||
normalized = value.strip()
|
||||
return normalized or default
|
||||
|
||||
|
||||
def ai_operator_default_config() -> AIOperatorConfig:
|
||||
agent_name = _env_field("AI_OPERATOR_AGENT_NAME", "Айнур")
|
||||
company_name = _env_field("AI_OPERATOR_COMPANY_NAME", "DigiOps")
|
||||
return AIOperatorConfig(
|
||||
agent_name="Айнур",
|
||||
company_name="DigiOps",
|
||||
base_system_prompt=(
|
||||
"Ты Айнур, единый ИИ-оператор контакт-центра DigiOps для звонков, Telegram и других каналов. "
|
||||
"Всегда сохраняй одну и ту же личность: тебя зовут Айнур. Если клиент спрашивает, кто ты или как "
|
||||
"тебя зовут, отвечай, что ты Айнур. Отвечай естественно, кратко и по делу. Когда говоришь о себе, "
|
||||
"используй женский род: могла, смогла, сделала, готова, проверила, нашла. Не завершай диалог "
|
||||
"самостоятельно и говори «до свидания» только если клиент явно попрощался или попросил завершить "
|
||||
"разговор. Телефонные номера читай по цифрам. Не используй Markdown, URL или таблицы."
|
||||
agent_name=agent_name,
|
||||
company_name=company_name,
|
||||
base_system_prompt=_env_field(
|
||||
"AI_OPERATOR_BASE_SYSTEM_PROMPT",
|
||||
(
|
||||
f"Ты {agent_name}, единый ИИ-оператор контакт-центра {company_name} для звонков, Telegram и других "
|
||||
f"каналов. Всегда сохраняй одну и ту же личность: тебя зовут {agent_name}. Если клиент спрашивает, "
|
||||
f"кто ты или как тебя зовут, отвечай, что ты {agent_name}. Отвечай естественно, кратко и по делу. "
|
||||
"Когда говоришь о себе, используй женский род: могла, смогла, сделала, готова, проверила, нашла. "
|
||||
"Не завершай диалог самостоятельно и говори «до свидания» только если клиент явно попрощался или "
|
||||
"попросил завершить разговор. Телефонные номера читай по цифрам. Не используй Markdown, URL или "
|
||||
"таблицы."
|
||||
),
|
||||
),
|
||||
identity_reply_ru=_env_field(
|
||||
"AI_OPERATOR_IDENTITY_REPLY_RU",
|
||||
f"Я {agent_name}, оператор контакт-центра {company_name}. Чем могу помочь?",
|
||||
),
|
||||
identity_reply_kz=_env_field(
|
||||
"AI_OPERATOR_IDENTITY_REPLY_KZ",
|
||||
f"Мен {agent_name}мын, {company_name} байланыс орталығының операторымын. Қалай көмектесе аламын?",
|
||||
),
|
||||
voice_greeting_ru=_env_field(
|
||||
"AI_OPERATOR_VOICE_GREETING_RU",
|
||||
f"Здравствуйте. Я {agent_name}. Подскажите, пожалуйста, чем помочь. (тест деплоя)",
|
||||
),
|
||||
voice_greeting_kz=_env_field(
|
||||
"AI_OPERATOR_VOICE_GREETING_KZ",
|
||||
f"Сәлеметсіз бе. Мен {agent_name}мын. Қалай көмектесе аламын?",
|
||||
),
|
||||
identity_reply_ru="Я Айнур, оператор контакт-центра DigiOps. Чем могу помочь?",
|
||||
identity_reply_kz="Мен Айнурмын, DigiOps байланыс орталығының операторымын. Қалай көмектесе аламын?",
|
||||
voice_greeting_ru="Здравствуйте. Я Айнур. Подскажите, пожалуйста, чем помочь. (тест деплоя)",
|
||||
voice_greeting_kz="Сәлеметсіз бе. Мен Айнурмын. Қалай көмектесе аламын?",
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,39 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Iterable
|
||||
|
||||
# Fixed, code-level conversational signals — not FAQ topics. These strings are
|
||||
# already used as intent literals across ai_orchestrator_service/app.py,
|
||||
# voice.py, and asserted directly in tests; centralized here rather than
|
||||
# renamed so every call site validates against the same set.
|
||||
CONTROL_INTENTS: frozenset[str] = frozenset(
|
||||
{
|
||||
"identity_question",
|
||||
"handoff_request",
|
||||
"sensitive_request",
|
||||
"resolution_confirmed",
|
||||
"clarification",
|
||||
"kb_answer",
|
||||
"unknown",
|
||||
}
|
||||
)
|
||||
|
||||
UNKNOWN_INTENT = "unknown"
|
||||
|
||||
|
||||
def normalize_intent(raw: str | None, *, known_topic_codes: Iterable[str] = ()) -> str:
|
||||
"""Validate a model-produced intent against control intents and KB topic codes.
|
||||
|
||||
`known_topic_codes` are the `intent_code` values of the KB articles actually
|
||||
shown to the model for this turn — anything else the model invents collapses
|
||||
to UNKNOWN_INTENT rather than being trusted verbatim.
|
||||
"""
|
||||
candidate = str(raw or "").strip()
|
||||
if not candidate:
|
||||
return UNKNOWN_INTENT
|
||||
if candidate in CONTROL_INTENTS:
|
||||
return candidate
|
||||
normalized_topic_codes = {str(code or "").strip().upper() for code in known_topic_codes if str(code or "").strip()}
|
||||
if candidate.upper() in normalized_topic_codes:
|
||||
return candidate.upper()
|
||||
return UNKNOWN_INTENT
|
||||
@@ -19,12 +19,41 @@ _MIN_SOFT_MATCH_LENGTH = 4
|
||||
_NON_WORD_RE = re.compile(r"[^\w]+", re.UNICODE)
|
||||
_SPACE_RE = re.compile(r"\s+")
|
||||
|
||||
# Common RU/KZ greetings, confirmations, pronouns, and particles that carry no
|
||||
# topical signal on their own. Without this, a caller utterance as thin as
|
||||
# "да" or "хорошо" could still exact-token-match some unrelated KB article
|
||||
# that happens to contain that word in its body, and get returned as the
|
||||
# top/only search result — read back to the caller as if it were the answer
|
||||
# to their question. Filtering these keeps _score_row's exact-token match
|
||||
# meaningful: a match now requires an actual content word.
|
||||
_STOPWORDS = frozenset(
|
||||
{
|
||||
# RU: greetings / confirmations / fillers
|
||||
"алло", "ага", "да", "неа", "нет", "ой", "ок", "окей", "угу", "ясно",
|
||||
"ладно", "хорошо", "понял", "поняла", "привет", "здравствуйте",
|
||||
"добрый", "день", "вечер", "утро", "слышу", "слышно", "спасибо",
|
||||
"пожалуйста", "извините", "простите", "алло",
|
||||
# RU: pronouns / conjunctions / particles with no topical content
|
||||
"я", "ты", "вы", "мы", "он", "она", "они", "это", "то", "и", "а",
|
||||
"но", "или", "что", "как", "где", "когда", "если", "чтобы", "для",
|
||||
"из", "по", "на", "в", "с", "у", "о", "же", "ли", "бы", "не", "ну",
|
||||
"вот", "просто", "есть", "быть", "можно", "нужно", "надо", "уже",
|
||||
# KZ: greetings / confirmations / fillers
|
||||
"иә", "ия", "жоқ", "жарайды", "түсінікті", "рахмет", "сәлем",
|
||||
"сәлеметсіз", "бе", "кешіріңіз",
|
||||
# KZ: pronouns / conjunctions / particles
|
||||
"мен", "сен", "сіз", "біз", "ол", "олар", "және", "бірақ", "немесе",
|
||||
"не", "қалай", "қайда", "қашан", "үшін", "туралы",
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
class KBSearchRow(Protocol):
|
||||
id: int
|
||||
title: str
|
||||
body: str
|
||||
tags_json: str
|
||||
intent_code: str | None
|
||||
|
||||
|
||||
T = TypeVar("T", bound=KBSearchRow)
|
||||
@@ -37,7 +66,11 @@ def normalize_kb_text(value: str | None) -> str:
|
||||
|
||||
|
||||
def tokenize_kb_text(value: str | None) -> list[str]:
|
||||
return [token for token in normalize_kb_text(value).split(" ") if len(token) >= _MIN_TOKEN_LENGTH]
|
||||
return [
|
||||
token
|
||||
for token in normalize_kb_text(value).split(" ")
|
||||
if len(token) >= _MIN_TOKEN_LENGTH and token not in _STOPWORDS
|
||||
]
|
||||
|
||||
|
||||
def search_kb_rows(
|
||||
|
||||
@@ -363,6 +363,12 @@ class VoiceEventIn(BaseModel):
|
||||
"recording.ready",
|
||||
"call.connected",
|
||||
"call.transferred",
|
||||
"AgentReserved",
|
||||
"AgentRinging",
|
||||
"AgentNoAnswer",
|
||||
"AgentConnected",
|
||||
"TransferCompleted",
|
||||
"TransferFailed",
|
||||
]
|
||||
call_id: str
|
||||
interaction_id: str | None = None
|
||||
@@ -983,6 +989,7 @@ class KBCategoryOut(KBCategoryCreate):
|
||||
class KBArticleCreate(BaseModel):
|
||||
category_id: str
|
||||
article_group_id: str | None = None
|
||||
intent_code: str | None = None
|
||||
language: str = "ru"
|
||||
title: str
|
||||
body: str
|
||||
@@ -991,6 +998,7 @@ class KBArticleCreate(BaseModel):
|
||||
|
||||
class KBArticleUpdate(BaseModel):
|
||||
article_group_id: str | None = None
|
||||
intent_code: str | None = None
|
||||
language: str | None = None
|
||||
title: str | None = None
|
||||
body: str | None = None
|
||||
@@ -1311,11 +1319,16 @@ class EscalationOut(BaseModel):
|
||||
summary: str | None = None
|
||||
status: str
|
||||
assigned_agent_id: str | None = None
|
||||
attempt_count: int = 0
|
||||
requested_at: str
|
||||
connected_at: str | None = None
|
||||
completed_at: str | None = None
|
||||
|
||||
|
||||
class RoutingAgentStatusIn(BaseModel):
|
||||
status: str = Field(min_length=1)
|
||||
|
||||
|
||||
class RoutingAgentReserveIn(BaseModel):
|
||||
call_id: str = Field(min_length=1)
|
||||
level: AgentLevel
|
||||
|
||||
@@ -583,6 +583,7 @@ def _apply_runtime_schema_compatibility() -> None:
|
||||
if "kb_articles" in table_names:
|
||||
columns = _table_columns(inspector, "kb_articles")
|
||||
_add_column_if_missing(conn, columns, "kb_articles", "article_group_id", "VARCHAR(64)")
|
||||
_add_column_if_missing(conn, columns, "kb_articles", "intent_code", "VARCHAR(64)")
|
||||
_add_column_if_missing(conn, columns, "kb_articles", "language", "VARCHAR(8) DEFAULT 'ru'")
|
||||
if "language" in columns:
|
||||
conn.execute(
|
||||
@@ -619,6 +620,13 @@ def _apply_runtime_schema_compatibility() -> None:
|
||||
"ON kb_articles(language)"
|
||||
)
|
||||
)
|
||||
if "idx_kb_articles_intent_code" not in indexes:
|
||||
conn.execute(
|
||||
text(
|
||||
"CREATE INDEX IF NOT EXISTS idx_kb_articles_intent_code "
|
||||
"ON kb_articles(intent_code)"
|
||||
)
|
||||
)
|
||||
|
||||
if "sales_automation_tasks" in table_names:
|
||||
columns = _table_columns(inspector, "sales_automation_tasks")
|
||||
|
||||
@@ -700,6 +700,7 @@ class KBArticleRow(Base):
|
||||
article_id: Mapped[str] = mapped_column(String(64), unique=True, index=True)
|
||||
category_id: Mapped[str] = mapped_column(String(64), index=True)
|
||||
article_group_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
|
||||
intent_code: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
|
||||
language: Mapped[str] = mapped_column(String(8), index=True, default="ru")
|
||||
title: Mapped[str] = mapped_column(String(512), index=True)
|
||||
body: Mapped[str] = mapped_column(Text)
|
||||
@@ -846,6 +847,9 @@ class EscalationRow(Base):
|
||||
summary: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||
status: Mapped[str] = mapped_column(String(32), index=True, default="requested")
|
||||
assigned_agent_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
|
||||
real_agent_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
|
||||
attempted_agent_ids_json: Mapped[str] = mapped_column(Text, default="[]")
|
||||
attempt_count: Mapped[int] = mapped_column(Integer, default=0)
|
||||
requested_at: Mapped[str] = mapped_column(String(64), index=True)
|
||||
connected_at: Mapped[str | None] = mapped_column(String(64), nullable=True)
|
||||
completed_at: Mapped[str | None] = mapped_column(String(64), nullable=True)
|
||||
|
||||
@@ -414,6 +414,7 @@ def seed_kb_article(
|
||||
*,
|
||||
language: str = "ru",
|
||||
article_group_id: str | None = None,
|
||||
intent_code: str | None = None,
|
||||
) -> dict[str, str]:
|
||||
session = get_session()
|
||||
try:
|
||||
@@ -434,6 +435,7 @@ def seed_kb_article(
|
||||
article_id=article_id,
|
||||
category_id=category_id,
|
||||
article_group_id=resolved_group_id,
|
||||
intent_code=intent_code,
|
||||
language=language,
|
||||
title=title,
|
||||
body=body,
|
||||
@@ -1347,6 +1349,88 @@ def test_voice_llm_guarded_decision_uses_operator_style_without_ai_or_kb(monkeyp
|
||||
assert "Do not say or imply that you are an AI" in system_prompt
|
||||
|
||||
|
||||
def test_voice_llm_decision_echoes_kb_article_intent_code(monkeypatch):
|
||||
monkeypatch.setenv("AI_VOICE_POLICY_MODE", "llm_guarded")
|
||||
monkeypatch.setattr(ai_module, "_ai_provider", lambda: "openai_compatible")
|
||||
|
||||
def _fake_structured(messages, **kwargs):
|
||||
del messages
|
||||
return {
|
||||
"language": "ru",
|
||||
"intent": "voucher_activation",
|
||||
"reply_text": "Подтвердите СМС с номера 1414 командой 21*1, затем завершите активацию в eGov.",
|
||||
"confidence": 0.9,
|
||||
"needs_handoff": False,
|
||||
"handoff_reason": None,
|
||||
"case_action": "keep_open",
|
||||
"kb_refs": ["kba_voucher_1"],
|
||||
"_model": "gpt-test",
|
||||
"_latency_ms": 30,
|
||||
"_finish_reason": "stop",
|
||||
}
|
||||
|
||||
monkeypatch.setattr(ai_module, "_request_structured_model_decision", _fake_structured)
|
||||
|
||||
kb_article = SimpleNamespace(
|
||||
article_id="kba_voucher_1",
|
||||
title="Активация ваучера",
|
||||
body="Подтвердите СМС 1414 командой 21*1, затем перейдите по ссылке и завершите в eGov Mobile.",
|
||||
intent_code="VOUCHER_ACTIVATION",
|
||||
)
|
||||
decision = voice_module._voice_decision(
|
||||
language="ru",
|
||||
customer=None,
|
||||
interaction=SimpleNamespace(interaction_id="int_voice_voucher", customer_id=None, status="new", queue_id="que_voice", subject="voucher"),
|
||||
transcript_text="Что делать с СМС от 1414?",
|
||||
transcript_window=[SimpleNamespace(speaker="caller", text="Что делать с СМС от 1414?", sequence_no=1, source_type="voice_asr", barge_in_interrupted=False, created_at=utc_now_iso())],
|
||||
kb_results=[kb_article],
|
||||
disclosure_required=False,
|
||||
)
|
||||
|
||||
assert decision["intent"] == "VOUCHER_ACTIVATION"
|
||||
|
||||
|
||||
def test_voice_llm_decision_rejects_invented_intent_not_in_kb_results(monkeypatch):
|
||||
monkeypatch.setenv("AI_VOICE_POLICY_MODE", "llm_guarded")
|
||||
monkeypatch.setattr(ai_module, "_ai_provider", lambda: "openai_compatible")
|
||||
|
||||
def _fake_structured(messages, **kwargs):
|
||||
del messages
|
||||
return {
|
||||
"language": "ru",
|
||||
"intent": "totally_made_up_intent",
|
||||
"reply_text": "Подтвердите СМС с номера 1414 командой 21*1.",
|
||||
"confidence": 0.9,
|
||||
"needs_handoff": False,
|
||||
"handoff_reason": None,
|
||||
"case_action": "keep_open",
|
||||
"kb_refs": ["kba_voucher_2"],
|
||||
"_model": "gpt-test",
|
||||
"_latency_ms": 30,
|
||||
"_finish_reason": "stop",
|
||||
}
|
||||
|
||||
monkeypatch.setattr(ai_module, "_request_structured_model_decision", _fake_structured)
|
||||
|
||||
kb_article = SimpleNamespace(
|
||||
article_id="kba_voucher_2",
|
||||
title="Активация ваучера",
|
||||
body="Подтвердите СМС 1414 командой 21*1.",
|
||||
intent_code="VOUCHER_ACTIVATION",
|
||||
)
|
||||
decision = voice_module._voice_decision(
|
||||
language="ru",
|
||||
customer=None,
|
||||
interaction=SimpleNamespace(interaction_id="int_voice_voucher_2", customer_id=None, status="new", queue_id="que_voice", subject="voucher"),
|
||||
transcript_text="Куда отправлять 21*1?",
|
||||
transcript_window=[SimpleNamespace(speaker="caller", text="Куда отправлять 21*1?", sequence_no=1, source_type="voice_asr", barge_in_interrupted=False, created_at=utc_now_iso())],
|
||||
kb_results=[kb_article],
|
||||
disclosure_required=False,
|
||||
)
|
||||
|
||||
assert decision["intent"] == "unknown"
|
||||
|
||||
|
||||
def test_voice_v2_fast_conversational_adds_ack_metadata_and_compacts_reply(monkeypatch):
|
||||
monkeypatch.setenv("AI_VOICE_POLICY_MODE", "v2_fast_conversational")
|
||||
monkeypatch.setattr(ai_module, "_ai_provider", lambda: "openai_compatible")
|
||||
@@ -1681,6 +1765,61 @@ def test_voice_llm_prompt_includes_context_summary_and_uses_12_segments():
|
||||
assert payload["history"][0]["sequence_no"] == 4
|
||||
|
||||
|
||||
def test_voice_llm_prompt_instructs_model_not_to_self_name_customer():
|
||||
messages = voice_module._voice_llm_prompt_messages(
|
||||
language="ru",
|
||||
customer=None,
|
||||
interaction=SimpleNamespace(
|
||||
interaction_id="int_voice_prompt_name",
|
||||
status="open",
|
||||
queue_id="que_voice",
|
||||
subject="schedule",
|
||||
customer_id="cus_voice_prompt_name",
|
||||
),
|
||||
transcript_text="Мне нужен график работы",
|
||||
transcript_window=[],
|
||||
conversation_summary_text="",
|
||||
kb_results=[],
|
||||
name_value="Ернур",
|
||||
name_status="name_obtained",
|
||||
)
|
||||
system_prompt = messages[0]["content"]
|
||||
assert "addressing the customer by name" in system_prompt
|
||||
|
||||
|
||||
def test_voice_reply_with_name_does_not_duplicate_inflected_name_form():
|
||||
# The model may address the customer using a grammatically declined form of
|
||||
# their name ("Данияре" instead of "Данияр"); an exact-token dedup check
|
||||
# would miss this and prepend the name a second time.
|
||||
reply = voice_module._voice_reply_with_name(
|
||||
"ru", "Здравствуйте, Данияре! Чем могу помочь?", "Данияр"
|
||||
)
|
||||
assert reply == "Здравствуйте, Данияре! Чем могу помочь?"
|
||||
|
||||
# A reply with no mention of the customer's name still gets it prefixed once.
|
||||
reply = voice_module._voice_reply_with_name("ru", "Чем могу помочь?", "Данияр")
|
||||
assert reply == "Данияр, Чем могу помочь?"
|
||||
|
||||
|
||||
def test_voice_reply_with_name_greet_mode_uses_one_of_two_fixed_forms():
|
||||
# Regular turns (name already known): just the name, never a greeting word.
|
||||
reply = voice_module._voice_reply_with_name("ru", "Чем могу помочь?", "Данияр", greet=False)
|
||||
assert reply == "Данияр, Чем могу помочь?"
|
||||
|
||||
# The turn the name is first learned: exactly "Здравствуйте, {name}, ...".
|
||||
reply = voice_module._voice_reply_with_name("ru", "Чем могу помочь?", "Данияр", greet=True)
|
||||
assert reply == "Здравствуйте, Данияр, Чем могу помочь?"
|
||||
|
||||
reply = voice_module._voice_reply_with_name("kz", "Немен көмектесе аламын?", "Ерлан", greet=True)
|
||||
assert reply == "Сәлеметсіз бе, Ерлан, Немен көмектесе аламын?"
|
||||
|
||||
# Still deduplicates even in greet mode if the model already named the customer.
|
||||
reply = voice_module._voice_reply_with_name(
|
||||
"ru", "Здравствуйте, Данияре! Чем могу помочь?", "Данияр", greet=True
|
||||
)
|
||||
assert reply == "Здравствуйте, Данияре! Чем могу помочь?"
|
||||
|
||||
|
||||
def test_voice_postprocess_reply_uses_summary_context_when_raw_window_lost_topic():
|
||||
reply_text = voice_module._voice_postprocess_reply_text(
|
||||
language="ru",
|
||||
|
||||
@@ -855,3 +855,31 @@ def test_push_voice_ai_telephony_event_call_ended_returns_detached_safe_payload(
|
||||
}
|
||||
assert closed == [("avs_runtime_call_ended", "call_ended")]
|
||||
assert orchestrator_calls == [("POST", "/ai/voice/sessions/avs_runtime_call_ended/close")]
|
||||
|
||||
|
||||
def test_media_registration_emotive_ack_rotation_is_not_gated_on_v2_queue_eligibility(monkeypatch):
|
||||
# AI_VOICE_V2_QUEUE_CODES defaults to "voice_lab_ai" only, so a queue outside
|
||||
# that allowlist runs the plain v1 ack path. Emotive-ack rotation and prebaked-ack
|
||||
# caching must still apply there — otherwise every filler collapses to the single
|
||||
# fixed "Секунду." fallback string instead of rotating, and gets synthesized live
|
||||
# via ElevenLabs each time instead of served from the pre-cached ack bank, adding
|
||||
# real TTS round-trip time to the one phrase whose job is to hide that latency.
|
||||
monkeypatch.delenv("AI_VOICE_V2_QUEUE_CODES", raising=False)
|
||||
monkeypatch.delenv("AI_VOICE_V2_EMOTIVE_ACK_ENABLED", raising=False)
|
||||
monkeypatch.delenv("AI_VOICE_V2_PREBAKED_ACK_ENABLED", raising=False)
|
||||
row = VoiceAISessionRow(
|
||||
session_id="avs_runtime_emotive_ack_v1",
|
||||
call_id="call_runtime_emotive_ack_v1",
|
||||
interaction_id="int_runtime_emotive_ack_v1",
|
||||
ai_session_id="ais_runtime_emotive_ack_v1",
|
||||
language="ru",
|
||||
media_uuid="media-emotive-ack-v1",
|
||||
queue_id="que_not_v2_eligible",
|
||||
agent_profile="voice_support",
|
||||
)
|
||||
|
||||
registration = runtime_module._media_registration_from_row(row, queue_code="queue_outside_v2_allowlist")
|
||||
|
||||
assert registration.voice_v2_enabled is False
|
||||
assert registration.voice_v2_emotive_ack is True
|
||||
assert registration.voice_v2_prebaked_ack is True
|
||||
|
||||
@@ -2195,6 +2195,7 @@ def test_disabled_bridge_startup_does_not_require_singleton_guard(monkeypatch):
|
||||
monkeypatch.delenv("ASTERISK_BRIDGE_ENABLED", raising=False)
|
||||
monkeypatch.setattr(bridge_module, "_ami_loop", _wait_until_stopped)
|
||||
monkeypatch.setattr(bridge_module, "_failed_retry_loop", _wait_until_stopped)
|
||||
monkeypatch.setattr(bridge_module, "_acw_sweep_loop", _wait_until_stopped)
|
||||
monkeypatch.setattr(
|
||||
bridge_module,
|
||||
"_try_acquire_bridge_singleton_guard",
|
||||
@@ -2202,7 +2203,7 @@ def test_disabled_bridge_startup_does_not_require_singleton_guard(monkeypatch):
|
||||
)
|
||||
|
||||
bridge_module._startup()
|
||||
assert len(bridge_module._background_threads_alive()) == 2
|
||||
assert len(bridge_module._background_threads_alive()) == 3
|
||||
|
||||
|
||||
def test_shutdown_releases_singleton_guard(monkeypatch):
|
||||
|
||||
@@ -0,0 +1,190 @@
|
||||
import json
|
||||
|
||||
import pytest
|
||||
|
||||
import services.asterisk_bridge_service.voice_ai as voice_ai
|
||||
from services.routing_service import engine as routing_engine
|
||||
from services.shared.core import new_id, utc_now_iso
|
||||
from services.shared.db import get_session
|
||||
from services.shared.sql_init import init_sql_schema
|
||||
from services.shared.sql_models import AgentRow, AsteriskCallLinkRow, EscalationRow
|
||||
|
||||
|
||||
def _make_agent(session, *, extension: str, status: str = "AVAILABLE"):
|
||||
now = utc_now_iso()
|
||||
row = AgentRow(
|
||||
agent_id=new_id("agt"),
|
||||
tenant_ids_json="[]",
|
||||
extension=extension,
|
||||
endpoint=None,
|
||||
display_name=f"Agent {extension}",
|
||||
level="L2",
|
||||
skills_json="[]",
|
||||
status=status,
|
||||
max_concurrent_calls=1,
|
||||
enabled=True,
|
||||
created_at=now,
|
||||
updated_at=now,
|
||||
)
|
||||
session.add(row)
|
||||
session.commit()
|
||||
session.refresh(row)
|
||||
return row
|
||||
|
||||
|
||||
def _make_call_link(session, *, call_id: str):
|
||||
now = utc_now_iso()
|
||||
link = AsteriskCallLinkRow(
|
||||
call_id=call_id,
|
||||
linked_id=call_id,
|
||||
queue_code="voice_lab_ai",
|
||||
queue_id="que_test",
|
||||
interaction_id="int_test",
|
||||
status="active",
|
||||
telephony_status="ringing",
|
||||
channel_name=f"PJSIP/tele2-kazgaz-{call_id}",
|
||||
current_level="L2",
|
||||
started_at=now,
|
||||
updated_at=now,
|
||||
)
|
||||
session.add(link)
|
||||
session.commit()
|
||||
session.refresh(link)
|
||||
return link
|
||||
|
||||
|
||||
def _make_ringing_escalation(session, *, call_id: str, agent: AgentRow):
|
||||
now = utc_now_iso()
|
||||
escalation = EscalationRow(
|
||||
escalation_id=new_id("esc"),
|
||||
call_id=call_id,
|
||||
tenant_id=None,
|
||||
from_level="L1",
|
||||
to_level="L2",
|
||||
reason_code="AI_UNABLE_TO_RESOLVE",
|
||||
required_skills_json="[]",
|
||||
priority=3,
|
||||
status="ringing",
|
||||
assigned_agent_id=agent.extension,
|
||||
real_agent_id=agent.agent_id,
|
||||
attempted_agent_ids_json=json.dumps([agent.agent_id]),
|
||||
requested_at=now,
|
||||
)
|
||||
session.add(escalation)
|
||||
session.commit()
|
||||
session.refresh(escalation)
|
||||
return escalation
|
||||
|
||||
|
||||
def _patch_routing_over_http(monkeypatch, session):
|
||||
"""Make voice_ai's HTTP-facing routing helpers operate on the same test session directly."""
|
||||
|
||||
def fake_reserve(*, call_id, level, tenant_id, required_skills=None, exclude_agent_ids=None):
|
||||
agent = routing_engine.reserve_agent(
|
||||
session,
|
||||
call_id=call_id,
|
||||
level=level,
|
||||
tenant_id=tenant_id,
|
||||
required_skills=required_skills,
|
||||
exclude_agent_ids=exclude_agent_ids,
|
||||
)
|
||||
if agent is None:
|
||||
return None
|
||||
return {
|
||||
"agent_id": agent.agent_id,
|
||||
"extension": agent.extension,
|
||||
"endpoint": agent.endpoint,
|
||||
"display_name": agent.display_name,
|
||||
}
|
||||
|
||||
def fake_release_by_agent_id(agent_id, *, next_status="AVAILABLE"):
|
||||
routing_engine.release_agent_by_id(session, agent_id=agent_id, next_status=next_status)
|
||||
|
||||
def fake_set_status(agent_id, status):
|
||||
routing_engine.set_agent_status(session, agent_id=agent_id, status=status)
|
||||
|
||||
monkeypatch.setattr(voice_ai, "_reserve_routing_agent", fake_reserve)
|
||||
monkeypatch.setattr(voice_ai, "routing_release_by_agent_id", fake_release_by_agent_id)
|
||||
monkeypatch.setattr(voice_ai, "set_routing_agent_status", fake_set_status)
|
||||
monkeypatch.setattr(voice_ai, "_append_escalation_timeline", lambda *a, **k: None)
|
||||
monkeypatch.setattr(voice_ai, "_emit_escalation_event", lambda *a, **k: None)
|
||||
monkeypatch.setattr(voice_ai, "_resolve_handoff_channel", lambda session, link: link.channel_name)
|
||||
|
||||
|
||||
def test_retry_escalation_no_answer_moves_to_next_available_agent(monkeypatch):
|
||||
init_sql_schema()
|
||||
redirected_to: list[str] = []
|
||||
|
||||
session = get_session()
|
||||
try:
|
||||
_patch_routing_over_http(monkeypatch, session)
|
||||
monkeypatch.setattr(
|
||||
voice_ai,
|
||||
"_redirect_channel_to_agent",
|
||||
lambda *, channel, extension: redirected_to.append(extension),
|
||||
)
|
||||
agent_a = _make_agent(session, extension="2001")
|
||||
agent_b = _make_agent(session, extension="2002")
|
||||
routing_engine.reserve_agent(session, call_id="call-retry-1", level="L2", tenant_id=None, required_skills=[])
|
||||
link = _make_call_link(session, call_id="call-retry-1")
|
||||
escalation = _make_ringing_escalation(session, call_id="call-retry-1", agent=agent_a)
|
||||
|
||||
voice_ai.retry_escalation_no_answer(session, call_id="call-retry-1", dial_outcome="NOANSWER")
|
||||
|
||||
session.refresh(escalation)
|
||||
session.refresh(agent_a)
|
||||
session.refresh(agent_b)
|
||||
|
||||
assert escalation.attempt_count == 1
|
||||
assert escalation.status == "ringing"
|
||||
assert escalation.real_agent_id == agent_b.agent_id
|
||||
assert escalation.assigned_agent_id == agent_b.extension
|
||||
assert json.loads(escalation.attempted_agent_ids_json) == [agent_a.agent_id, agent_b.agent_id]
|
||||
|
||||
assert agent_a.status == "AVAILABLE"
|
||||
assert agent_a.current_call_id is None
|
||||
assert agent_b.status == "RINGING"
|
||||
assert agent_b.current_call_id == "call-retry-1"
|
||||
|
||||
assert redirected_to == [agent_b.extension]
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
|
||||
def test_retry_escalation_no_answer_exhausts_pool_marks_failed(monkeypatch):
|
||||
init_sql_schema()
|
||||
|
||||
session = get_session()
|
||||
try:
|
||||
_patch_routing_over_http(monkeypatch, session)
|
||||
monkeypatch.setattr(voice_ai, "_redirect_channel_to_agent", lambda **kwargs: None)
|
||||
agent_a = _make_agent(session, extension="3001")
|
||||
routing_engine.reserve_agent(session, call_id="call-retry-2", level="L2", tenant_id=None, required_skills=[])
|
||||
link = _make_call_link(session, call_id="call-retry-2")
|
||||
escalation = _make_ringing_escalation(session, call_id="call-retry-2", agent=agent_a)
|
||||
|
||||
voice_ai.retry_escalation_no_answer(session, call_id="call-retry-2", dial_outcome="NOANSWER")
|
||||
|
||||
session.refresh(escalation)
|
||||
session.refresh(agent_a)
|
||||
|
||||
assert escalation.attempt_count == 1
|
||||
assert escalation.status == "failed"
|
||||
assert agent_a.status == "AVAILABLE"
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
|
||||
def test_retry_escalation_no_answer_ignores_calls_without_ringing_escalation(monkeypatch):
|
||||
init_sql_schema()
|
||||
called = []
|
||||
|
||||
session = get_session()
|
||||
try:
|
||||
_patch_routing_over_http(monkeypatch, session)
|
||||
monkeypatch.setattr(voice_ai, "_redirect_channel_to_agent", lambda **kwargs: called.append(kwargs))
|
||||
_make_call_link(session, call_id="call-no-escalation")
|
||||
voice_ai.retry_escalation_no_answer(session, call_id="call-no-escalation", dial_outcome="NOANSWER")
|
||||
assert called == []
|
||||
finally:
|
||||
session.close()
|
||||
@@ -0,0 +1,34 @@
|
||||
from services.shared.intents import CONTROL_INTENTS, UNKNOWN_INTENT, normalize_intent
|
||||
|
||||
|
||||
def test_control_intent_passes_through_unchanged():
|
||||
assert normalize_intent("handoff_request") == "handoff_request"
|
||||
assert normalize_intent("kb_answer", known_topic_codes=["VOUCHER_ACTIVATION"]) == "kb_answer"
|
||||
|
||||
|
||||
def test_matching_topic_code_passes_through_case_insensitively():
|
||||
assert normalize_intent("voucher_activation", known_topic_codes=["VOUCHER_ACTIVATION"]) == "VOUCHER_ACTIVATION"
|
||||
assert normalize_intent(" Voucher_Activation ", known_topic_codes=["voucher_activation"]) == "VOUCHER_ACTIVATION"
|
||||
|
||||
|
||||
def test_unknown_topic_code_falls_back_to_unknown():
|
||||
assert normalize_intent("made_up_intent", known_topic_codes=["VOUCHER_ACTIVATION"]) == UNKNOWN_INTENT
|
||||
assert normalize_intent("voucher_activation", known_topic_codes=[]) == UNKNOWN_INTENT
|
||||
|
||||
|
||||
def test_empty_or_missing_intent_falls_back_to_unknown():
|
||||
assert normalize_intent(None) == UNKNOWN_INTENT
|
||||
assert normalize_intent("") == UNKNOWN_INTENT
|
||||
assert normalize_intent(" ") == UNKNOWN_INTENT
|
||||
|
||||
|
||||
def test_control_intents_frozenset_matches_documented_values():
|
||||
assert CONTROL_INTENTS == {
|
||||
"identity_question",
|
||||
"handoff_request",
|
||||
"sensitive_request",
|
||||
"resolution_confirmed",
|
||||
"clarification",
|
||||
"kb_answer",
|
||||
"unknown",
|
||||
}
|
||||
@@ -114,6 +114,61 @@ def test_search_kb_rows_prefers_title_and_tags_over_body_only_mentions():
|
||||
assert results[0].title == rows[0].title
|
||||
|
||||
|
||||
def test_search_kb_rows_ignores_stopword_only_query():
|
||||
# A caller confirming the language ("да") should never surface an
|
||||
# unrelated FAQ article just because that article's body happens to
|
||||
# contain the word "да" somewhere in ordinary prose.
|
||||
rows = [
|
||||
_row(
|
||||
row_id=1,
|
||||
title="Активация ваучера",
|
||||
body="Да, подтвердите СМС с номера 1414 командой 21*1.",
|
||||
tags=["ваучер"],
|
||||
),
|
||||
]
|
||||
|
||||
assert search_kb_rows(rows, "да", limit=5) == []
|
||||
assert search_kb_rows(rows, "Хорошо", limit=5) == []
|
||||
|
||||
|
||||
def test_search_kb_rows_still_matches_real_content_word_amid_fillers():
|
||||
rows = [
|
||||
_row(
|
||||
row_id=1,
|
||||
title="Активация ваучера",
|
||||
body="Подтвердите СМС с номера 1414 командой 21*1.",
|
||||
tags=["ваучер"],
|
||||
),
|
||||
_row(
|
||||
row_id=2,
|
||||
title="График работы",
|
||||
body="Филиалы работают с 9 до 18.",
|
||||
tags=["график"],
|
||||
),
|
||||
]
|
||||
|
||||
results = search_kb_rows(rows, "да, у меня вопрос про ваучер", limit=5)
|
||||
|
||||
assert results
|
||||
assert results[0].title == "Активация ваучера"
|
||||
|
||||
|
||||
def test_search_kb_rows_single_real_word_query_still_matches():
|
||||
rows = [
|
||||
_row(
|
||||
row_id=1,
|
||||
title="Активация ваучера",
|
||||
body="Подтвердите СМС с номера 1414 командой 21*1.",
|
||||
tags=["ваучер"],
|
||||
),
|
||||
]
|
||||
|
||||
results = search_kb_rows(rows, "ваучер", limit=5)
|
||||
|
||||
assert results
|
||||
assert results[0].title == "Активация ваучера"
|
||||
|
||||
|
||||
def test_search_kb_rows_breaks_ties_by_newer_id():
|
||||
older = _row(
|
||||
row_id=10,
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
from datetime import datetime, timedelta, timezone
|
||||
import json
|
||||
|
||||
from sqlalchemy import select
|
||||
@@ -134,3 +135,50 @@ def test_reserve_agent_excludes_disabled_and_excluded_ids():
|
||||
assert reserved.agent_id != excluded.agent_id
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
|
||||
def test_set_agent_status_transitions_ringing_to_talking():
|
||||
init_sql_schema()
|
||||
session = get_session()
|
||||
try:
|
||||
agent = _make_agent(session, level="L2")
|
||||
routing_engine.reserve_agent(session, call_id="call-status", level="L2", tenant_id=None, required_skills=[])
|
||||
|
||||
ringing = routing_engine.set_agent_status(session, agent_id=agent.agent_id, status="RINGING")
|
||||
assert ringing is not None
|
||||
assert ringing.status == "RINGING"
|
||||
|
||||
talking = routing_engine.set_agent_status(session, agent_id=agent.agent_id, status="TALKING")
|
||||
assert talking.status == "TALKING"
|
||||
|
||||
missing = routing_engine.set_agent_status(session, agent_id="unknown-agent", status="AVAILABLE")
|
||||
assert missing is None
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
|
||||
def test_sweep_after_call_work_releases_only_expired_agents():
|
||||
init_sql_schema()
|
||||
session = get_session()
|
||||
try:
|
||||
stale = _make_agent(session, level="L2", status="AFTER_CALL_WORK")
|
||||
fresh = _make_agent(session, level="L2", status="AFTER_CALL_WORK")
|
||||
|
||||
stale_ts = (datetime.now(timezone.utc) - timedelta(seconds=120)).replace(microsecond=0).isoformat()
|
||||
stale.updated_at = stale_ts
|
||||
stale.current_call_id = "call-stale"
|
||||
session.commit()
|
||||
|
||||
swept = routing_engine.sweep_after_call_work(session, older_than_seconds=30)
|
||||
swept_ids = {a.agent_id for a in swept}
|
||||
|
||||
assert stale.agent_id in swept_ids
|
||||
assert fresh.agent_id not in swept_ids
|
||||
|
||||
session.refresh(stale)
|
||||
session.refresh(fresh)
|
||||
assert stale.status == "AVAILABLE"
|
||||
assert stale.current_call_id is None
|
||||
assert fresh.status == "AFTER_CALL_WORK"
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
@@ -40,6 +40,46 @@ def test_kb_lite_search():
|
||||
assert len(search.json()) >= 1
|
||||
|
||||
|
||||
def test_kb_article_intent_code_round_trips_through_create_get_update():
|
||||
client = TestClient(kb_app)
|
||||
headers = {"X-User": "analyst", "X-Role": "analyst"}
|
||||
|
||||
cat = client.post(
|
||||
"/knowledge/categories",
|
||||
json={"name": "Vouchers", "description": "Voucher help"},
|
||||
headers=headers,
|
||||
)
|
||||
assert cat.status_code == 200
|
||||
category_id = cat.json()["category_id"]
|
||||
|
||||
created = client.post(
|
||||
"/knowledge/articles",
|
||||
json={
|
||||
"category_id": category_id,
|
||||
"title": "Активация ваучера",
|
||||
"body": "Подтвердите СМС 1414 командой 21*1.",
|
||||
"tags": ["voucher"],
|
||||
"intent_code": "voucher_activation",
|
||||
},
|
||||
headers=headers,
|
||||
)
|
||||
assert created.status_code == 200
|
||||
assert created.json()["intent_code"] == "VOUCHER_ACTIVATION"
|
||||
article_id = created.json()["article_id"]
|
||||
|
||||
fetched = client.get(f"/knowledge/articles/{article_id}")
|
||||
assert fetched.status_code == 200
|
||||
assert fetched.json()["intent_code"] == "VOUCHER_ACTIVATION"
|
||||
|
||||
updated = client.patch(
|
||||
f"/knowledge/articles/{article_id}",
|
||||
json={"intent_code": "voucher_activation_v2"},
|
||||
headers=headers,
|
||||
)
|
||||
assert updated.status_code == 200
|
||||
assert updated.json()["intent_code"] == "VOUCHER_ACTIVATION_V2"
|
||||
|
||||
|
||||
def test_kb_lite_search_ranks_title_over_body_only_matches():
|
||||
client = TestClient(kb_app)
|
||||
headers = {"X-User": "analyst", "X-Role": "analyst"}
|
||||
|
||||
Reference in New Issue
Block a user