mirror of
https://github.com/dartdavros/chatballs.git
synced 2026-10-05 01:14:58 +03:00
✨ feat(conversations): ответ AI считается отдельно от приёма входящих
Ход AI выполнялся прямо в приёме сообщения: цикл опроса мессенджеров и HTTP-запрос виджета ждали провайдера, держа открытой транзакцию организации. Один ход — это два обращения к модели (эмбеддинг и чат) по тридцать секунд с двумя повторами, то есть до трёх минут, и всё это время ни одно входящее по всей установке не забиралось. Владелец видел это как «бот залипает»: сайт и MAX на одном агенте отвечали с задержками или молчали. Приём теперь доводит дело до записи сообщения и ставит событие `conversation.ai_turn_requested`. Ход считает роль событий воркера (`run_worker --role=events`) короткими транзакциями, между которыми остаются походы к провайдеру и в мессенджер. Туда же уехала расшифровка голосовых — последнее обращение наружу из цикла опроса. Воркер разделён на роли: `poller` опрашивает подключения и ведёт периодические работы (один экземпляр — курсоры и паузы после сбоя живут в его памяти), `events` разбирает outbox и масштабируется репликами (`CHATBALLS_EVENT_WORKERS`, по умолчанию две). Роль `all` осталась для разработки. Очередь событий научилась двум вещам: события одного диалога не выдаются параллельно (иначе два ответа приезжают клиенту вперемешку) и событие, взятое упавшим процессом, возвращается в очередь по истечении аренды. Попутно убраны мины, которые тот же залип и продлевали: - ход клиенту ограничен своим таймаутом (CHATBALLS_AI_TURN_TIMEOUT, 20 с) и сроком годности (CHATBALLS_AI_TURN_DEADLINE_SECONDS, 120 с) — просроченный ход не зовёт модель, а передаёт диалог оператору; - отказ провайдера по существу запроса (4xx, кроме 429) больше не повторяется трижды по таймауту; - предохранитель провайдера считает сбои по ключу «организация + интеграция», а не один на процесс: отозванный ключ одной организации гасил AI у всех; - потолок паузы после сбоя опроса — минута вместо четверти часа: он был компромиссом ради журнала однопоточного воркера. Виджет узнаёт, что ответ считается, по признаку `thinking` в ленте. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
1 parent
c8b8246b35
commit
83b1119d6b
37 files changed
+1702
-430
No files matched your search
@@ -1,4 +1,20 @@
|
||||
"""Обращения к LLM-провайдеру: подготовка, сам вызов и запись в журнал.
|
||||
|
||||
Вызов провайдера ждёт ответа десятки секунд, а подготовка и журнал — это
|
||||
обращения к базе. В одной функции они означают открытую транзакцию на всё время
|
||||
ожидания, а вместе с ней занятое соединение из пула и RLS-контекст
|
||||
(chatballs.tenancy.middleware). Поэтому шаги разделены: `prepare_*` и `record_*`
|
||||
вызывают внутри транзакции, `run_*` — вне её.
|
||||
|
||||
`invoke_chat` и `embed_texts` остаются для мест, где ждать под транзакцией не
|
||||
жалко: индексация знаний, предпросмотр карточки агента, тесты. Ход диалога с
|
||||
клиентом ходит по шагам (chatballs.ai.turn).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from dataclasses import dataclass
|
||||
|
||||
from django.conf import settings
|
||||
|
||||
@@ -15,19 +31,129 @@ from chatballs.ai.provider.base import (
|
||||
from chatballs.ai.provider.factory import get_provider
|
||||
from chatballs.ai.provider.resilience import CircuitBreaker, call_with_resilience
|
||||
|
||||
_breaker = CircuitBreaker()
|
||||
# Предохранитель считает сбои по ключу «организация + интеграция»: провайдер у
|
||||
# каждой организации свой, и отозванный ключ одной не имеет отношения к AI
|
||||
# остальных. Общий на процесс предохранитель гасил AI у всех сразу.
|
||||
_breakers: dict[tuple[int, int], CircuitBreaker] = {}
|
||||
|
||||
|
||||
def _prepare_invocation(*, channel, requested_model: str | None) -> tuple[LLMProvider, str]:
|
||||
def _breaker(key: tuple[int, int]) -> CircuitBreaker:
|
||||
breaker = _breakers.get(key)
|
||||
if breaker is None:
|
||||
breaker = CircuitBreaker()
|
||||
_breakers[key] = breaker
|
||||
return breaker
|
||||
|
||||
|
||||
def reset_breakers() -> None:
|
||||
"""Для тестов: забыть накопленные сбои провайдеров."""
|
||||
|
||||
_breakers.clear()
|
||||
|
||||
|
||||
def _breaker_key(channel) -> tuple[int, int]:
|
||||
"""Ключ предохранителя. Без канала провайдер может быть только тестовым —
|
||||
считать сбои там не по чему, и общий ключ (0, 0) никому не мешает."""
|
||||
|
||||
if channel is None:
|
||||
return (0, 0)
|
||||
return (channel.organization_id, routing.integration_id(channel))
|
||||
|
||||
|
||||
def _elapsed_ms(started: float) -> int:
|
||||
return int((time.monotonic() - started) * 1000)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class ChatJob:
|
||||
"""Всё для похода к модели, уже прочитанное из базы."""
|
||||
|
||||
provider: LLMProvider
|
||||
model: str
|
||||
messages: list[ChatMessage]
|
||||
breaker_key: tuple[int, int]
|
||||
params: dict | None = None
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class EmbeddingJob:
|
||||
"""То же для эмбеддингов: вектор считается тем же провайдером организации."""
|
||||
|
||||
provider: LLMProvider
|
||||
model: str
|
||||
texts: list[str]
|
||||
breaker_key: tuple[int, int]
|
||||
|
||||
|
||||
def _effective_model(channel, requested_model: str | None) -> str:
|
||||
# BYOK — единственный режим (ADR-CHATBALLS-0042 §3): модель берётся из интеграции
|
||||
# организации с fallback на модель агента. Без интеграции модель остаётся
|
||||
# агентской: тестовый провайдер работает, прод упадёт в get_provider штатно.
|
||||
agent = channel.ai_agent
|
||||
agent = getattr(channel, "ai_agent", None)
|
||||
fallback = str(getattr(agent, "model", "") or "")
|
||||
if requested_model:
|
||||
return requested_model
|
||||
try:
|
||||
effective_model = routing.resolve_model(channel, fallback_model=agent.model)
|
||||
return routing.resolve_model(channel, fallback_model=fallback)
|
||||
except routing.IntegrationNotConfigured:
|
||||
effective_model = agent.model
|
||||
return get_provider(channel=channel), requested_model or effective_model
|
||||
return fallback
|
||||
|
||||
|
||||
def prepare_chat(
|
||||
*,
|
||||
channel,
|
||||
messages: list[ChatMessage],
|
||||
model: str | None = None,
|
||||
params: dict | None = None,
|
||||
timeout: float | None = None,
|
||||
) -> ChatJob:
|
||||
"""Шаг в транзакции: провайдер, модель и очищенный от ПДн текст запроса."""
|
||||
|
||||
return ChatJob(
|
||||
provider=get_provider(channel=channel, timeout=timeout),
|
||||
model=_effective_model(channel, model),
|
||||
messages=[ChatMessage(role=item.role, content=redact(item.content)) for item in messages],
|
||||
breaker_key=_breaker_key(channel),
|
||||
params=params,
|
||||
)
|
||||
|
||||
|
||||
def run_chat(job: ChatJob) -> ChatResult:
|
||||
"""Шаг без транзакции: обращение к провайдеру."""
|
||||
|
||||
return call_with_resilience(
|
||||
lambda: job.provider.chat(messages=job.messages, model=job.model, params=job.params),
|
||||
retries=settings.CHATBALLS_AI_MAX_RETRIES,
|
||||
breaker=_breaker(job.breaker_key),
|
||||
)
|
||||
|
||||
|
||||
def record_chat(
|
||||
*,
|
||||
channel,
|
||||
job: ChatJob,
|
||||
purpose: str,
|
||||
result: ChatResult | None = None,
|
||||
error: Exception | None = None,
|
||||
latency_ms: int = 0,
|
||||
used_fragment_ids: list | None = None,
|
||||
) -> None:
|
||||
"""Шаг в транзакции: строка журнала вызовов — и об успехе, и об отказе."""
|
||||
|
||||
LlmInvocation.objects.create(
|
||||
organization=channel.organization,
|
||||
channel=channel,
|
||||
purpose=purpose,
|
||||
operation="chat",
|
||||
model=result.model if result is not None else job.model,
|
||||
prompt_tokens=result.prompt_tokens if result else 0,
|
||||
completion_tokens=result.completion_tokens if result else 0,
|
||||
total_tokens=result.total_tokens if result else 0,
|
||||
latency_ms=latency_ms,
|
||||
status=LlmInvocationStatus.SUCCESS if result is not None else LlmInvocationStatus.ERROR,
|
||||
error="" if error is None else str(error)[:1000],
|
||||
used_fragment_ids=used_fragment_ids or [],
|
||||
)
|
||||
|
||||
|
||||
def invoke_chat(
|
||||
@@ -39,44 +165,82 @@ def invoke_chat(
|
||||
params: dict | None = None,
|
||||
used_fragment_ids: list | None = None,
|
||||
) -> ChatResult:
|
||||
provider, model = _prepare_invocation(channel=channel, requested_model=model)
|
||||
"""Три шага подряд, в транзакции вызывающего: там, где ждать не жалко."""
|
||||
|
||||
safe_messages = [ChatMessage(role=item.role, content=redact(item.content)) for item in messages]
|
||||
job = prepare_chat(channel=channel, messages=messages, model=model, params=params)
|
||||
started = time.monotonic()
|
||||
try:
|
||||
result: ChatResult = call_with_resilience(
|
||||
lambda: provider.chat(messages=safe_messages, model=model, params=params),
|
||||
retries=settings.CHATBALLS_AI_MAX_RETRIES,
|
||||
breaker=_breaker,
|
||||
)
|
||||
result = run_chat(job)
|
||||
except ProviderError as error:
|
||||
LlmInvocation.objects.create(
|
||||
organization=channel.organization,
|
||||
record_chat(
|
||||
channel=channel,
|
||||
job=job,
|
||||
purpose=purpose,
|
||||
operation="chat",
|
||||
model=model,
|
||||
status=LlmInvocationStatus.ERROR,
|
||||
error=str(error)[:1000],
|
||||
latency_ms=int((time.monotonic() - started) * 1000),
|
||||
error=error,
|
||||
latency_ms=_elapsed_ms(started),
|
||||
)
|
||||
raise
|
||||
record_chat(
|
||||
channel=channel,
|
||||
job=job,
|
||||
purpose=purpose,
|
||||
result=result,
|
||||
latency_ms=_elapsed_ms(started),
|
||||
used_fragment_ids=used_fragment_ids,
|
||||
)
|
||||
return result
|
||||
|
||||
return_result = result
|
||||
|
||||
def prepare_embedding(
|
||||
*,
|
||||
channel,
|
||||
texts: list[str],
|
||||
model: str,
|
||||
timeout: float | None = None,
|
||||
) -> EmbeddingJob:
|
||||
"""Шаг в транзакции: провайдер эмбеддингов организации."""
|
||||
|
||||
return EmbeddingJob(
|
||||
provider=get_provider(channel=channel, timeout=timeout),
|
||||
model=model,
|
||||
texts=texts,
|
||||
breaker_key=_breaker_key(channel),
|
||||
)
|
||||
|
||||
|
||||
def run_embedding(job: EmbeddingJob) -> list[EmbeddingResult]:
|
||||
"""Шаг без транзакции: обращение к провайдеру."""
|
||||
|
||||
return call_with_resilience(
|
||||
lambda: job.provider.embed(texts=job.texts, model=job.model),
|
||||
retries=settings.CHATBALLS_AI_MAX_RETRIES,
|
||||
breaker=_breaker(job.breaker_key),
|
||||
)
|
||||
|
||||
|
||||
def record_embedding(
|
||||
*,
|
||||
channel=None,
|
||||
organization=None,
|
||||
model: str,
|
||||
purpose: str,
|
||||
results: list[EmbeddingResult],
|
||||
latency_ms: int = 0,
|
||||
) -> None:
|
||||
"""Шаг в транзакции: строка журнала."""
|
||||
|
||||
tokens = sum(result.tokens for result in results)
|
||||
LlmInvocation.objects.create(
|
||||
organization=channel.organization,
|
||||
organization=channel.organization if channel else organization,
|
||||
channel=channel,
|
||||
purpose=purpose,
|
||||
operation="chat",
|
||||
model=result.model,
|
||||
prompt_tokens=result.prompt_tokens,
|
||||
completion_tokens=result.completion_tokens,
|
||||
total_tokens=result.total_tokens,
|
||||
latency_ms=int((time.monotonic() - started) * 1000),
|
||||
operation="embedding",
|
||||
model=model,
|
||||
prompt_tokens=tokens,
|
||||
total_tokens=tokens,
|
||||
latency_ms=latency_ms,
|
||||
status=LlmInvocationStatus.SUCCESS,
|
||||
used_fragment_ids=used_fragment_ids or [],
|
||||
)
|
||||
return return_result
|
||||
|
||||
|
||||
def embed_texts(
|
||||
@@ -87,21 +251,17 @@ def embed_texts(
|
||||
model: str,
|
||||
purpose: str = "retrieval",
|
||||
) -> list[EmbeddingResult]:
|
||||
provider = get_provider(channel=channel)
|
||||
results: list[EmbeddingResult] = call_with_resilience(
|
||||
lambda: provider.embed(texts=texts, model=model),
|
||||
retries=settings.CHATBALLS_AI_MAX_RETRIES,
|
||||
breaker=_breaker,
|
||||
)
|
||||
tokens = sum(result.tokens for result in results)
|
||||
LlmInvocation.objects.create(
|
||||
organization=channel.organization if channel else organization,
|
||||
"""Три шага подряд: индексация знаний и прочие неинтерактивные места."""
|
||||
|
||||
job = prepare_embedding(channel=channel, texts=texts, model=model)
|
||||
started = time.monotonic()
|
||||
results = run_embedding(job)
|
||||
record_embedding(
|
||||
channel=channel,
|
||||
purpose=purpose,
|
||||
operation="embedding",
|
||||
organization=organization,
|
||||
model=model,
|
||||
prompt_tokens=tokens,
|
||||
total_tokens=tokens,
|
||||
status=LlmInvocationStatus.SUCCESS,
|
||||
purpose=purpose,
|
||||
results=results,
|
||||
latency_ms=_elapsed_ms(started),
|
||||
)
|
||||
return results
|
||||
@@ -35,6 +35,16 @@ class ProviderError(Exception):
|
||||
"""Transient/technical provider failure (eligible for retry / circuit breaker)."""
|
||||
|
||||
|
||||
class ProviderRejected(ProviderError):
|
||||
"""Отказ, который повтором не лечится: провайдер не принял сам запрос.
|
||||
|
||||
Неверный ключ, несуществующая модель, слишком длинный контекст. Повтор
|
||||
потратит ещё один таймаут и получит тот же ответ, а клиент всё это время
|
||||
ждёт ответа. «Слишком часто» (429) сюда не относится — это как раз тот
|
||||
случай, когда повторить стоит.
|
||||
"""
|
||||
|
||||
|
||||
class LLMProvider(abc.ABC):
|
||||
name: str = "base"
|
||||
|
||||
|
||||
@@ -13,7 +13,7 @@ def _test_provider() -> LLMProvider:
|
||||
return LocalProvider()
|
||||
|
||||
|
||||
def get_provider(*, channel=None) -> LLMProvider:
|
||||
def get_provider(*, channel=None, timeout: float | None = None) -> LLMProvider:
|
||||
"""Resolve the organization's own provider (BYOK, ADR-CHATBALLS-0042 §3).
|
||||
|
||||
The test adapter is an explicit test-surface override. Managed platform
|
||||
@@ -29,10 +29,10 @@ def get_provider(*, channel=None) -> LLMProvider:
|
||||
raise ProviderError(
|
||||
t("ai.provider_not_configured")
|
||||
)
|
||||
return routing.resolve_provider(channel)
|
||||
return routing.resolve_provider(channel, timeout=timeout)
|
||||
|
||||
|
||||
def get_transcription_provider(*, channel=None) -> LLMProvider:
|
||||
def get_transcription_provider(*, channel=None, timeout: float | None = None) -> LLMProvider:
|
||||
"""Провайдер расшифровки голосовых.
|
||||
|
||||
Отличается от `get_provider` одним: агент может расшифровывать другим
|
||||
@@ -42,4 +42,4 @@ def get_transcription_provider(*, channel=None) -> LLMProvider:
|
||||
return _test_provider()
|
||||
if channel is None:
|
||||
raise ProviderError(t("ai.provider_not_configured"))
|
||||
return routing.resolve_transcription_provider(channel)
|
||||
return routing.resolve_transcription_provider(channel, timeout=timeout)
|
||||
@@ -37,7 +37,19 @@ import json
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
|
||||
from chatballs.ai.provider.base import ChatMessage, ChatResult, EmbeddingResult, ProviderError
|
||||
from chatballs.ai.provider.base import (
|
||||
|
||||
ChatMessage,
|
||||
|
||||
ChatResult,
|
||||
|
||||
EmbeddingResult,
|
||||
|
||||
ProviderError,
|
||||
|
||||
ProviderRejected,
|
||||
|
||||
)
|
||||
from chatballs.i18n import t
|
||||
from chatballs.integrations.proxy import build_opener
|
||||
|
||||
@@ -72,6 +84,20 @@ def post_json(*, base_url: str, path: str, api_key: str, payload: dict, timeout:
|
||||
|
||||
return json.loads(response.read().decode("utf-8"))
|
||||
|
||||
# Отказ самого провайдера разбирается отдельно: 4xx (кроме 429) — это ключ,
|
||||
|
||||
# модель или размер запроса, и повтор даст тот же ответ через ещё один таймаут.
|
||||
|
||||
except urllib.error.HTTPError as error:
|
||||
|
||||
detail = error.read().decode("utf-8", "replace")[:300]
|
||||
|
||||
if error.code != 429 and 400 <= error.code < 500:
|
||||
|
||||
raise ProviderRejected(f"HTTP {error.code}: {detail}") from error
|
||||
|
||||
raise ProviderError(f"HTTP {error.code}: {detail}") from error
|
||||
|
||||
# http.client.HTTPException covers IncompleteRead/BadStatusLine (dropped reply)
|
||||
|
||||
# — those are not OSError, so they would slip past ProviderError otherwise.
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import time
|
||||
from collections.abc import Callable
|
||||
|
||||
from chatballs.ai.provider.base import ProviderError
|
||||
from chatballs.ai.provider.base import ProviderError, ProviderRejected
|
||||
|
||||
|
||||
class CircuitBreakerOpen(ProviderError):
|
||||
@@ -44,6 +44,10 @@ def call_with_resilience(
|
||||
breaker.before()
|
||||
try:
|
||||
result = func()
|
||||
except ProviderRejected:
|
||||
# Провайдер отказал по существу запроса: повторять нечего, и
|
||||
# предохранитель тут ни при чём — сам провайдер жив и отвечает.
|
||||
raise
|
||||
except ProviderError:
|
||||
if breaker is not None:
|
||||
breaker.on_failure()
|
||||
|
||||
@@ -40,10 +40,14 @@ class IntegrationNotConfigured(ProviderError):
|
||||
"""
|
||||
|
||||
|
||||
def resolve_provider(channel) -> LLMProvider:
|
||||
"""Build the BYOK LLMProvider from the channel's explicit integration."""
|
||||
def resolve_provider(channel, *, timeout: float | None = None) -> LLMProvider:
|
||||
"""Build the BYOK LLMProvider from the channel's explicit integration.
|
||||
|
||||
`timeout` переопределяет срок ожидания ответа: интерактивному ходу диалога
|
||||
отведено меньше, чем индексации знаний (chatballs.ai.turn).
|
||||
"""
|
||||
integration = _channel_integration(channel)
|
||||
return _provider_from_integration(integration)
|
||||
return _provider_from_integration(integration, timeout=timeout)
|
||||
|
||||
|
||||
def resolve_provider_and_model(channel, *, fallback_model: str) -> tuple[LLMProvider, str]:
|
||||
@@ -90,9 +94,9 @@ def _transcription_integration(channel) -> Integration:
|
||||
return integration
|
||||
|
||||
|
||||
def resolve_transcription_provider(channel) -> LLMProvider:
|
||||
def resolve_transcription_provider(channel, *, timeout: float | None = None) -> LLMProvider:
|
||||
"""Провайдер расшифровки: отдельная интеграция агента либо провайдер ответов."""
|
||||
return _provider_from_integration(_transcription_integration(channel))
|
||||
return _provider_from_integration(_transcription_integration(channel), timeout=timeout)
|
||||
|
||||
|
||||
def resolve_transcription_model(channel) -> str:
|
||||
@@ -120,14 +124,29 @@ def _channel_integration(channel) -> Integration:
|
||||
return integration
|
||||
|
||||
|
||||
def _provider_from_integration(integration: Integration) -> LLMProvider:
|
||||
def integration_id(channel) -> int:
|
||||
"""Идентификатор интеграции канала; 0 — интеграции нет.
|
||||
|
||||
Нужен там, где интеграция — ключ, а не источник настроек: предохранитель
|
||||
считает сбои по конкретному ключу организации (chatballs.ai.invocation).
|
||||
"""
|
||||
try:
|
||||
return _channel_integration(channel).id
|
||||
except IntegrationNotConfigured:
|
||||
return 0
|
||||
|
||||
|
||||
def _provider_from_integration(
|
||||
integration: Integration, *, timeout: float | None = None
|
||||
) -> LLMProvider:
|
||||
from django.conf import settings
|
||||
|
||||
wait = timeout or settings.CHATBALLS_AI_REQUEST_TIMEOUT
|
||||
if integration.provider == IntegrationProvider.OPENROUTER:
|
||||
return OpenRouterProvider(
|
||||
api_key=integration.secret,
|
||||
base_url=integration.config.get("base_url") or settings.CHATBALLS_OPENROUTER_BASE_URL,
|
||||
timeout=settings.CHATBALLS_AI_REQUEST_TIMEOUT,
|
||||
timeout=wait,
|
||||
proxy_url=integration.config.get("proxy_url", ""),
|
||||
)
|
||||
if integration.provider == IntegrationProvider.DEMO:
|
||||
@@ -136,7 +155,7 @@ def _provider_from_integration(integration: Integration) -> LLMProvider:
|
||||
return CustomProvider(
|
||||
api_key=integration.secret,
|
||||
base_url=integration.config["base_url"],
|
||||
timeout=settings.CHATBALLS_AI_REQUEST_TIMEOUT,
|
||||
timeout=wait,
|
||||
proxy_url=integration.config.get("proxy_url", ""),
|
||||
)
|
||||
raise IntegrationNotConfigured(
|
||||
|
||||
@@ -43,6 +43,26 @@ def semantic_search(
|
||||
)
|
||||
|
||||
|
||||
def merge_hits(
|
||||
agent: AIAgent,
|
||||
query: str,
|
||||
query_vector: list[float] | None,
|
||||
*,
|
||||
limit: int = 5,
|
||||
) -> list[KnowledgeFragment]:
|
||||
"""Оба поиска и их склейка — шаг в транзакции, без обращений наружу.
|
||||
|
||||
Вектор считается отдельно (chatballs.ai.turn): поход за эмбеддингом — это
|
||||
сеть, и держать ради него транзакцию незачем. Без вектора остаётся
|
||||
лексический поиск: знания находятся хуже, но находятся.
|
||||
"""
|
||||
semantic = semantic_search(agent, query_vector, limit=limit) if query_vector else []
|
||||
lexical = lexical_search(agent, query, limit=limit)
|
||||
seen = {fragment.id for fragment in semantic}
|
||||
merged = semantic + [fragment for fragment in lexical if fragment.id not in seen]
|
||||
return merged[:limit]
|
||||
|
||||
|
||||
class KnowledgeRetriever:
|
||||
"""Hybrid retriever: semantic (pgvector) primary, lexical (Postgres FTS) complementary."""
|
||||
|
||||
@@ -56,10 +76,6 @@ class KnowledgeRetriever:
|
||||
model=settings.CHATBALLS_AI_EMBEDDING_MODEL,
|
||||
purpose="retrieval_query",
|
||||
)[0].vector
|
||||
semantic = semantic_search(agent, query_vector, limit=limit)
|
||||
except ProviderError:
|
||||
semantic = []
|
||||
lexical = lexical_search(agent, query, limit=limit)
|
||||
seen = {fragment.id for fragment in semantic}
|
||||
merged = semantic + [fragment for fragment in lexical if fragment.id not in seen]
|
||||
return merged[:limit]
|
||||
query_vector = None
|
||||
return merge_hits(agent, query, query_vector, limit=limit)
|
||||
@@ -119,15 +119,19 @@ def knowledge_catalog(agent: AIAgent) -> str:
|
||||
return "\n\n".join(parts)
|
||||
|
||||
|
||||
def run_agent_turn(
|
||||
def build_turn_messages(
|
||||
*,
|
||||
agent: AIAgent,
|
||||
message: str,
|
||||
history: list[dict] | None = None,
|
||||
fragments: list[KnowledgeFragment],
|
||||
style_guard: bool = True,
|
||||
) -> AgentTurnResult:
|
||||
fragments = KnowledgeRetriever().retrieve(agent=agent, query=message, limit=5)
|
||||
) -> list[ChatMessage]:
|
||||
"""Промпт хода целиком: инструкции агента, каталог знаний, найденное, история.
|
||||
|
||||
Только чтение базы и склейка строк — обращений наружу здесь нет, поэтому
|
||||
сборку можно держать внутри транзакции (chatballs.ai.turn).
|
||||
"""
|
||||
messages: list[ChatMessage] = []
|
||||
system_prompt = agent_system_prompt(agent)
|
||||
if system_prompt:
|
||||
@@ -156,7 +160,30 @@ def run_agent_turn(
|
||||
ChatMessage(role=str(item.get("role", "user")), content=str(item.get("content", "")))
|
||||
)
|
||||
messages.append(ChatMessage(role="user", content=message))
|
||||
return messages
|
||||
|
||||
|
||||
def run_agent_turn(
|
||||
*,
|
||||
agent: AIAgent,
|
||||
message: str,
|
||||
history: list[dict] | None = None,
|
||||
style_guard: bool = True,
|
||||
) -> AgentTurnResult:
|
||||
"""Ход агента целиком, в транзакции вызывающего.
|
||||
|
||||
Остаётся для мест, где ждать провайдера под транзакцией не жалко:
|
||||
предпросмотр на карточке агента и тесты. Ход диалога с клиентом идёт
|
||||
шагами, вне транзакции (chatballs.ai.turn).
|
||||
"""
|
||||
fragments = KnowledgeRetriever().retrieve(agent=agent, query=message, limit=5)
|
||||
messages = build_turn_messages(
|
||||
agent=agent,
|
||||
message=message,
|
||||
history=history,
|
||||
fragments=fragments,
|
||||
style_guard=style_guard,
|
||||
)
|
||||
result = invoke_chat(
|
||||
channel=agent.channel,
|
||||
messages=messages,
|
||||
|
||||
@@ -0,0 +1,175 @@
|
||||
"""Ход агента по шагам: транзакция — сеть — транзакция — сеть — транзакция.
|
||||
|
||||
Ответ клиенту складывается из двух обращений к провайдеру (вектор вопроса и
|
||||
сам ответ) и нескольких обращений к базе между ними. Сделанные подряд, они
|
||||
держат транзакцию организации открытой всё время ожидания провайдера — а это
|
||||
минуты (chatballs.ai.invocation). Здесь работа разложена так, чтобы каждое
|
||||
обращение к базе шло своей короткой транзакцией, а походы наружу оставались
|
||||
между ними.
|
||||
|
||||
Порядок шагов у вызывающего (chatballs.conversations.ai_turn):
|
||||
|
||||
1. в транзакции: `plan_query_embedding`
|
||||
2. вне транзакции: `run_query_embedding`
|
||||
3. в транзакции: `plan_chat`
|
||||
4. вне транзакции: `run_turn_chat`
|
||||
5. в транзакции: `record_turn` и запись ответа
|
||||
|
||||
Шаги `run_*` ошибок провайдера не поднимают: отказ — это такой же результат
|
||||
хода, его пишут в журнал и разбирают в диалоге (передачей оператору).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from dataclasses import dataclass, field
|
||||
|
||||
from django.conf import settings
|
||||
|
||||
from chatballs.ai.invocation import (
|
||||
ChatJob,
|
||||
EmbeddingJob,
|
||||
prepare_chat,
|
||||
prepare_embedding,
|
||||
record_chat,
|
||||
record_embedding,
|
||||
run_chat,
|
||||
run_embedding,
|
||||
)
|
||||
from chatballs.ai.models import AIAgent
|
||||
from chatballs.ai.provider.base import ChatResult, EmbeddingResult, ProviderError
|
||||
from chatballs.ai.retrieval import merge_hits
|
||||
from chatballs.ai.runtime import build_turn_messages
|
||||
|
||||
FRAGMENT_LIMIT = 5
|
||||
|
||||
|
||||
def _elapsed_ms(started: float) -> int:
|
||||
return int((time.monotonic() - started) * 1000)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class QueryEmbedding:
|
||||
"""Вектор вопроса. Пустой вектор — обычное дело: остаётся лексический поиск."""
|
||||
|
||||
vector: list[float] | None = None
|
||||
model: str = ""
|
||||
latency_ms: int = 0
|
||||
results: list[EmbeddingResult] = field(default_factory=list)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class TurnPlan:
|
||||
"""Готовый запрос к модели и то, на чём он основан."""
|
||||
|
||||
job: ChatJob
|
||||
fragment_ids: list[int]
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class TurnAnswer:
|
||||
"""Итог похода к модели: либо ответ, либо отказ, и сколько это заняло."""
|
||||
|
||||
result: ChatResult | None = None
|
||||
error: ProviderError | None = None
|
||||
latency_ms: int = 0
|
||||
|
||||
|
||||
def plan_query_embedding(*, agent: AIAgent, query: str) -> EmbeddingJob | None:
|
||||
"""Шаг в транзакции: чем считать вектор вопроса. None — считать нечем."""
|
||||
|
||||
if not query.strip():
|
||||
return None
|
||||
try:
|
||||
return prepare_embedding(
|
||||
channel=agent.channel,
|
||||
texts=[query],
|
||||
model=settings.CHATBALLS_AI_EMBEDDING_MODEL,
|
||||
timeout=settings.CHATBALLS_AI_TURN_TIMEOUT,
|
||||
)
|
||||
except ProviderError:
|
||||
# Провайдера нет или он не настроен: семантический поиск необязателен.
|
||||
return None
|
||||
|
||||
|
||||
def run_query_embedding(job: EmbeddingJob | None) -> QueryEmbedding:
|
||||
"""Шаг без транзакции: обращение к провайдеру за вектором."""
|
||||
|
||||
if job is None:
|
||||
return QueryEmbedding()
|
||||
started = time.monotonic()
|
||||
try:
|
||||
results = run_embedding(job)
|
||||
except ProviderError:
|
||||
# Знания найдутся лексическим поиском; ход из-за этого не срывается.
|
||||
return QueryEmbedding(latency_ms=_elapsed_ms(started))
|
||||
return QueryEmbedding(
|
||||
vector=results[0].vector if results else None,
|
||||
model=job.model,
|
||||
latency_ms=_elapsed_ms(started),
|
||||
results=results,
|
||||
)
|
||||
|
||||
|
||||
def plan_chat(
|
||||
*,
|
||||
agent: AIAgent,
|
||||
message: str,
|
||||
history: list[dict] | None = None,
|
||||
embedding: QueryEmbedding | None = None,
|
||||
style_guard: bool = True,
|
||||
) -> TurnPlan:
|
||||
"""Шаг в транзакции: поиск знаний, сборка промпта и выбор модели.
|
||||
|
||||
Заодно здесь оседает журнальная строка о векторе вопроса: считали его
|
||||
снаружи транзакции, а писать её всё равно в базу.
|
||||
"""
|
||||
embedding = embedding or QueryEmbedding()
|
||||
if embedding.results:
|
||||
record_embedding(
|
||||
channel=agent.channel,
|
||||
model=embedding.model,
|
||||
purpose="retrieval_query",
|
||||
results=embedding.results,
|
||||
latency_ms=embedding.latency_ms,
|
||||
)
|
||||
fragments = merge_hits(agent, message, embedding.vector, limit=FRAGMENT_LIMIT)
|
||||
job = prepare_chat(
|
||||
channel=agent.channel,
|
||||
messages=build_turn_messages(
|
||||
agent=agent,
|
||||
message=message,
|
||||
history=history,
|
||||
fragments=fragments,
|
||||
style_guard=style_guard,
|
||||
),
|
||||
model=agent.model,
|
||||
params=agent.model_params or None,
|
||||
timeout=settings.CHATBALLS_AI_TURN_TIMEOUT,
|
||||
)
|
||||
return TurnPlan(job=job, fragment_ids=[fragment.id for fragment in fragments])
|
||||
|
||||
|
||||
def run_turn_chat(plan: TurnPlan) -> TurnAnswer:
|
||||
"""Шаг без транзакции: обращение к модели за ответом."""
|
||||
|
||||
started = time.monotonic()
|
||||
try:
|
||||
result = run_chat(plan.job)
|
||||
except ProviderError as error:
|
||||
return TurnAnswer(error=error, latency_ms=_elapsed_ms(started))
|
||||
return TurnAnswer(result=result, latency_ms=_elapsed_ms(started))
|
||||
|
||||
|
||||
def record_turn(*, agent: AIAgent, plan: TurnPlan, answer: TurnAnswer) -> None:
|
||||
"""Шаг в транзакции: строка журнала вызовов — и об ответе, и об отказе."""
|
||||
|
||||
record_chat(
|
||||
channel=agent.channel,
|
||||
job=plan.job,
|
||||
purpose="agent_chat",
|
||||
result=answer.result,
|
||||
error=answer.error,
|
||||
latency_ms=answer.latency_ms,
|
||||
used_fragment_ids=plan.fragment_ids,
|
||||
)
|
||||
@@ -0,0 +1,302 @@
|
||||
"""Ход AI по входящему сообщению — отдельная работа, а не часть приёма.
|
||||
|
||||
Раньше ответ считался прямо в приёме: цикл опроса мессенджеров и HTTP-запрос
|
||||
виджета ждали провайдера минутами, держа открытой транзакцию организации, — и
|
||||
всё это время ни одно другое входящее не забиралось. Теперь приём доводит дело
|
||||
до записи сообщения и ставит ход в очередь событий, а считает его роль событий
|
||||
(`run_worker --role=events`), которую можно держать в нескольких процессах.
|
||||
|
||||
Границы транзакций здесь и есть главное: каждое обращение к базе идёт своей
|
||||
короткой транзакцией, походы к провайдеру и в мессенджер остаются между ними.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from dataclasses import dataclass
|
||||
from datetime import timedelta
|
||||
|
||||
from django.conf import settings
|
||||
from django.utils import timezone
|
||||
|
||||
from chatballs.ai.models import HISTORY_LIMIT_DEFAULT, AIAgent
|
||||
from chatballs.ai.provider.base import ProviderError
|
||||
from chatballs.ai.turn import (
|
||||
plan_chat,
|
||||
plan_query_embedding,
|
||||
record_turn,
|
||||
run_query_embedding,
|
||||
run_turn_chat,
|
||||
)
|
||||
from chatballs.conversations import ai_turn_result, transports
|
||||
from chatballs.conversations.models import (
|
||||
AiTurnState,
|
||||
ControlMode,
|
||||
Conversation,
|
||||
Message,
|
||||
MessageAuthor,
|
||||
MessageKind,
|
||||
)
|
||||
from chatballs.conversations.transcription import (
|
||||
TranscriptionJob,
|
||||
mark_transcription_failed,
|
||||
prepare_transcription,
|
||||
run_transcription,
|
||||
store_transcription,
|
||||
)
|
||||
from chatballs.events.services import DomainEvent, enqueue_event
|
||||
from chatballs.tenancy.context import TenantContext
|
||||
from chatballs.tenancy.database import tenant_atomic
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
AI_TURN_REQUESTED = "conversation.ai_turn_requested"
|
||||
# Агрегат события — диалог: ходы одного диалога обрабатываются строго по
|
||||
# очереди (chatballs.events.services.claim_next_outbox_event).
|
||||
AGGREGATE_TYPE = "Conversation"
|
||||
|
||||
_ROLE = {
|
||||
MessageAuthor.CONTACT: "user",
|
||||
MessageAuthor.AI: "assistant",
|
||||
MessageAuthor.OPERATOR: "assistant",
|
||||
MessageAuthor.SYSTEM: "system",
|
||||
}
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class Turn:
|
||||
"""Всё о ходе, прочитанное из базы первым шагом."""
|
||||
|
||||
message: Message
|
||||
conversation: Conversation
|
||||
agent: AIAgent
|
||||
user_id: str
|
||||
query: str
|
||||
history: list[dict]
|
||||
is_new_conversation: bool = False
|
||||
transcription_job: TranscriptionJob | None = None
|
||||
embedding_job: object | None = None
|
||||
# Ход прерван на подготовке, и клиенту есть что сказать: текст уходит ему
|
||||
# уже вне транзакции, как и обычный ответ.
|
||||
stopped: bool = False
|
||||
outgoing: str = ""
|
||||
|
||||
|
||||
def request_ai_turn(
|
||||
*,
|
||||
message: Message,
|
||||
user_id: str,
|
||||
context: TenantContext,
|
||||
is_new_conversation: bool = False,
|
||||
) -> None:
|
||||
"""Шаг в транзакции приёма: пометить сообщение и поставить ход в очередь."""
|
||||
|
||||
message.ai_turn_state = AiTurnState.PENDING
|
||||
message.save(update_fields=["ai_turn_state"])
|
||||
enqueue_event(
|
||||
DomainEvent(
|
||||
aggregate_type=AGGREGATE_TYPE,
|
||||
aggregate_id=str(message.conversation_id),
|
||||
event_type=AI_TURN_REQUESTED,
|
||||
payload={
|
||||
"messageId": message.id,
|
||||
"userId": user_id,
|
||||
# Про новый диалог операторов уже позвали при приёме: второй
|
||||
# оклик из-за нерасшифрованного голосового был бы лишним.
|
||||
"isNewConversation": is_new_conversation,
|
||||
},
|
||||
tenant_context=context,
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
def conversation_is_thinking(conversation_id: int) -> bool:
|
||||
"""Есть ли по диалогу ход, который прямо сейчас считается.
|
||||
|
||||
По этому же признаку виджет показывает клиенту, что ответ пишется.
|
||||
"""
|
||||
return Message.objects.filter(
|
||||
conversation_id=conversation_id,
|
||||
ai_turn_state__in=(AiTurnState.PENDING, AiTurnState.RUNNING),
|
||||
).exists()
|
||||
|
||||
|
||||
def _history(conversation: Conversation, limit: int) -> list[dict]:
|
||||
# С конца и с ограничением в базе: длинный диалог не поднимается в память
|
||||
# целиком ради последних сообщений. Самое новое — входящее, по которому
|
||||
# идёт ход, оно уходит модели отдельно.
|
||||
latest = conversation.messages.order_by("-created_at", "-id")[: limit + 1]
|
||||
prior = list(reversed(latest))[:-1]
|
||||
# Голосовые попадают в контекст стенограммой.
|
||||
return [
|
||||
{"role": _ROLE.get(m.author_type, "user"), "content": m.text or m.transcript}
|
||||
for m in prior
|
||||
if m.text or m.transcript
|
||||
]
|
||||
|
||||
|
||||
def _expired(message: Message) -> bool:
|
||||
deadline = timedelta(seconds=settings.CHATBALLS_AI_TURN_DEADLINE_SECONDS)
|
||||
return timezone.now() - message.created_at > deadline
|
||||
|
||||
|
||||
def _plan_transcription(message: Message, channel) -> TranscriptionJob | None:
|
||||
"""Голосовое без стенограммы: чем её снять. None — снимать нечем."""
|
||||
|
||||
if message.kind != MessageKind.VOICE or message.transcript:
|
||||
return None
|
||||
try:
|
||||
return prepare_transcription(channel, message)
|
||||
except ProviderError as error:
|
||||
logger.info("Voice transcription unavailable for message %s: %s", message.id, error)
|
||||
return None
|
||||
|
||||
|
||||
def _begin(*, message_id: int, user_id: str, is_new: bool, context: TenantContext) -> Turn | None:
|
||||
"""Шаг в транзакции: взять ход в работу — или отказаться от него.
|
||||
|
||||
Отказ здесь нормален и молчалив: событие могло приехать вторым заходом
|
||||
после сбоя, диалог мог уйти оператору, а ход мог пролежать в очереди
|
||||
дольше, чем ответ имеет смысл.
|
||||
"""
|
||||
message = (
|
||||
Message.objects.select_related(
|
||||
"conversation__channel__ai_agent",
|
||||
"conversation__channel__organization",
|
||||
"conversation__contact",
|
||||
"conversation__connection",
|
||||
)
|
||||
.filter(id=message_id)
|
||||
.first()
|
||||
)
|
||||
if message is None:
|
||||
return None
|
||||
if message.ai_turn_state not in (AiTurnState.PENDING, AiTurnState.RUNNING):
|
||||
return None
|
||||
conversation = message.conversation
|
||||
channel = conversation.channel
|
||||
agent = getattr(channel, "ai_agent", None)
|
||||
if conversation.control_mode != ControlMode.AI or agent is None or not agent.is_active:
|
||||
# Диалог успел уйти человеку либо агента отключили: отвечать не нужно.
|
||||
message.ai_turn_state = AiTurnState.DONE
|
||||
message.save(update_fields=["ai_turn_state"])
|
||||
return None
|
||||
turn = Turn(
|
||||
message=message,
|
||||
conversation=conversation,
|
||||
agent=agent,
|
||||
user_id=user_id,
|
||||
query=message.text or message.transcript,
|
||||
history=_history(conversation, agent.history_limit or HISTORY_LIMIT_DEFAULT),
|
||||
is_new_conversation=is_new,
|
||||
)
|
||||
message.ai_turn_state = AiTurnState.RUNNING
|
||||
message.save(update_fields=["ai_turn_state"])
|
||||
if _expired(message):
|
||||
turn.stopped = True
|
||||
turn.outgoing = ai_turn_result.store_failure(
|
||||
turn=turn, context=context, error="turn deadline passed"
|
||||
)
|
||||
return turn
|
||||
turn.transcription_job = _plan_transcription(message, channel)
|
||||
if turn.transcription_job is not None:
|
||||
# Вопрос станет известен после расшифровки — вместе с ним и вектор.
|
||||
return turn
|
||||
if not turn.query.strip():
|
||||
# Голосовое, которое нечем расшифровать, и прочее «отвечать не на что».
|
||||
ai_turn_result.store_voice_without_transcript(turn=turn, context=context)
|
||||
return None
|
||||
turn.embedding_job = plan_query_embedding(agent=agent, query=turn.query)
|
||||
return turn
|
||||
|
||||
|
||||
def _run_transcription(turn: Turn) -> str:
|
||||
"""Шаг без транзакции: голос в текст."""
|
||||
|
||||
try:
|
||||
return run_transcription(turn.transcription_job)
|
||||
except ProviderError as error:
|
||||
logger.info(
|
||||
"Voice transcription unavailable for message %s: %s", turn.message.id, error
|
||||
)
|
||||
return ""
|
||||
|
||||
|
||||
def _apply_transcript(*, turn: Turn, transcript: str, context: TenantContext) -> bool:
|
||||
"""Шаг в транзакции: сохранить стенограмму. False — хода не будет."""
|
||||
|
||||
if not transcript.strip():
|
||||
mark_transcription_failed(turn.message)
|
||||
ai_turn_result.store_voice_without_transcript(turn=turn, context=context)
|
||||
return False
|
||||
store_transcription(turn.message, transcript)
|
||||
turn.query = transcript
|
||||
turn.embedding_job = plan_query_embedding(agent=turn.agent, query=transcript)
|
||||
return True
|
||||
|
||||
|
||||
def _deliver(turn: Turn, text: str) -> None:
|
||||
"""Шаг без транзакции: ответ уходит клиенту в его канал.
|
||||
|
||||
Веб-виджет забирает ответ поллингом — для него отправка пустая.
|
||||
"""
|
||||
if not text or turn.conversation.connection is None:
|
||||
return
|
||||
transports.send_reply(
|
||||
turn.conversation.connection,
|
||||
chat_id=turn.conversation.external_chat_id,
|
||||
user_id=turn.user_id,
|
||||
text=text,
|
||||
)
|
||||
|
||||
|
||||
def run_requested_turn(payload: dict, context: TenantContext) -> None:
|
||||
"""Ход целиком: короткие транзакции и походы наружу между ними."""
|
||||
|
||||
message_id = int(payload.get("messageId") or 0)
|
||||
user_id = str(payload.get("userId") or "")
|
||||
is_new = bool(payload.get("isNewConversation"))
|
||||
with tenant_atomic(context):
|
||||
turn = _begin(message_id=message_id, user_id=user_id, is_new=is_new, context=context)
|
||||
if turn is None:
|
||||
return
|
||||
if turn.stopped:
|
||||
_deliver(turn, turn.outgoing)
|
||||
return
|
||||
|
||||
if turn.transcription_job is not None:
|
||||
transcript = _run_transcription(turn)
|
||||
with tenant_atomic(context):
|
||||
if not _apply_transcript(turn=turn, transcript=transcript, context=context):
|
||||
return
|
||||
|
||||
embedding = run_query_embedding(turn.embedding_job)
|
||||
failure = None
|
||||
with tenant_atomic(context):
|
||||
try:
|
||||
plan = plan_chat(
|
||||
agent=turn.agent,
|
||||
message=turn.query,
|
||||
history=turn.history,
|
||||
embedding=embedding,
|
||||
)
|
||||
except ProviderError as error:
|
||||
# Провайдер не настроен вовсе — тот же отказ хода, что и молчание
|
||||
# модели: клиент получает понятный текст, диалог уходит человеку.
|
||||
failure = ai_turn_result.store_failure(turn=turn, context=context, error=error)
|
||||
if failure is not None:
|
||||
_deliver(turn, failure)
|
||||
return
|
||||
|
||||
answer = run_turn_chat(plan)
|
||||
with tenant_atomic(context):
|
||||
record_turn(agent=turn.agent, plan=plan, answer=answer)
|
||||
if answer.error is not None:
|
||||
outgoing = ai_turn_result.store_failure(
|
||||
turn=turn, context=context, error=answer.error
|
||||
)
|
||||
else:
|
||||
outgoing = ai_turn_result.store_answer(
|
||||
turn=turn, context=context, text=answer.result.text
|
||||
)
|
||||
_deliver(turn, outgoing)
|
||||
@@ -0,0 +1,175 @@
|
||||
"""Что делать с результатом хода AI: ответ клиенту либо передача оператору.
|
||||
|
||||
Отделено от оркестрации (chatballs.conversations.ai_turn) намеренно: там —
|
||||
порядок шагов и границы транзакций, здесь — правила диалога. Обе функции
|
||||
вызывают внутри транзакции и обе возвращают текст, который нужно отправить
|
||||
клиенту: сама отправка — это сеть, и её место снаружи транзакции.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from django.utils import timezone
|
||||
|
||||
from chatballs.ai.runtime import HANDOFF_TOKEN
|
||||
from chatballs.conversations.models import (
|
||||
AiTurnState,
|
||||
ExpectedResponder,
|
||||
Message,
|
||||
MessageAuthor,
|
||||
SystemEvent,
|
||||
)
|
||||
from chatballs.conversations.queue import QUEUE_FIELDS, enter_queue
|
||||
from chatballs.i18n import customer_language, t
|
||||
from chatballs.notifications.models import NotificationAudience, NotificationType
|
||||
from chatballs.notifications.services import notify, notify_management
|
||||
from chatballs.tenancy.context import TenantContext
|
||||
|
||||
if TYPE_CHECKING: # pragma: no cover - только для подсказок типов
|
||||
from chatballs.conversations.ai_turn import Turn
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _finish(message: Message, state: str) -> None:
|
||||
message.ai_turn_state = state
|
||||
message.save(update_fields=["ai_turn_state"])
|
||||
|
||||
|
||||
def _contact_name(turn: Turn) -> str:
|
||||
return turn.conversation.contact.name or t("conversations.guest")
|
||||
|
||||
|
||||
def store_answer(*, turn: Turn, context: TenantContext, text: str) -> str:
|
||||
"""Ответ модели: запись в диалог и, если модель попросила, передача оператору.
|
||||
|
||||
Возвращает текст для отправки клиенту.
|
||||
"""
|
||||
conversation = turn.conversation
|
||||
reply = text
|
||||
handoff = HANDOFF_TOKEN in reply
|
||||
if handoff:
|
||||
reply = reply.replace(HANDOFF_TOKEN, "").strip()
|
||||
|
||||
Message.objects.create(conversation=conversation, author_type=MessageAuthor.AI, text=reply)
|
||||
conversation.last_activity_at = timezone.now()
|
||||
if handoff:
|
||||
enter_queue(conversation)
|
||||
else:
|
||||
conversation.expected_responder = ExpectedResponder.CUSTOMER
|
||||
conversation.save(update_fields=[*QUEUE_FIELDS, "last_activity_at"])
|
||||
_finish(turn.message, AiTurnState.DONE)
|
||||
|
||||
if handoff:
|
||||
Message.objects.create(
|
||||
conversation=conversation,
|
||||
author_type=MessageAuthor.SYSTEM,
|
||||
system_event=SystemEvent.AI_HANDED_OVER,
|
||||
text="AI передал диалог оператору",
|
||||
)
|
||||
notify(
|
||||
context=context,
|
||||
type=NotificationType.OPERATOR_REQUESTED,
|
||||
audience=NotificationAudience.OPERATORS,
|
||||
audience_group=conversation.group,
|
||||
title=f"AI передал диалог · {_contact_name(turn)}",
|
||||
title_key="notifications.ai_handed_over",
|
||||
text_params={"contact": _contact_name(turn)},
|
||||
body=turn.query[:120],
|
||||
target_id=conversation.id,
|
||||
source_type="Conversation",
|
||||
source_id=conversation.id,
|
||||
dedup_key=f"handoff:{conversation.id}",
|
||||
)
|
||||
return reply
|
||||
|
||||
|
||||
def store_voice_without_transcript(*, turn: Turn, context: TenantContext) -> None:
|
||||
"""Отвечать не на что: голосовое без стенограммы уходит оператору.
|
||||
|
||||
Это не сбой AI, и клиент не должен видеть извинений за поломку: ему просто
|
||||
ответит человек.
|
||||
"""
|
||||
conversation = turn.conversation
|
||||
enter_queue(conversation)
|
||||
conversation.save(update_fields=QUEUE_FIELDS)
|
||||
_finish(turn.message, AiTurnState.FAILED)
|
||||
if turn.is_new_conversation:
|
||||
# Про новый диалог операторов уже позвали при приёме.
|
||||
return
|
||||
notify(
|
||||
context=context,
|
||||
type=NotificationType.OPERATOR_REQUESTED,
|
||||
audience=NotificationAudience.OPERATORS,
|
||||
audience_group=conversation.group,
|
||||
title=f"Нужен оператор · {_contact_name(turn)}",
|
||||
title_key="notifications.operator_needed",
|
||||
text_params={"contact": _contact_name(turn)},
|
||||
body="Голосовое без расшифровки",
|
||||
body_key="notifications.voice_without_transcript",
|
||||
target_id=conversation.id,
|
||||
source_type="Conversation",
|
||||
source_id=conversation.id,
|
||||
dedup_key=f"media:{conversation.id}",
|
||||
)
|
||||
|
||||
|
||||
def store_failure(*, turn: Turn, context: TenantContext, error: object) -> str:
|
||||
"""Ответа не будет: диалог уходит оператору, клиент получает понятный текст.
|
||||
|
||||
Сбой AI не должен «терять» сообщение — ни отказ провайдера, ни ход,
|
||||
просроченный в очереди.
|
||||
"""
|
||||
conversation = turn.conversation
|
||||
channel = conversation.channel
|
||||
logger.warning("AI turn failed for conversation %s: %s", conversation.id, error)
|
||||
|
||||
enter_queue(conversation)
|
||||
conversation.last_activity_at = timezone.now()
|
||||
conversation.save(update_fields=[*QUEUE_FIELDS, "last_activity_at"])
|
||||
Message.objects.create(
|
||||
conversation=conversation,
|
||||
author_type=MessageAuthor.SYSTEM,
|
||||
system_event=SystemEvent.AI_UNAVAILABLE,
|
||||
text="AI недоступен — диалог передан оператору",
|
||||
)
|
||||
fallback = t(
|
||||
"conversations.ai_unavailable_reply",
|
||||
language=customer_language(channel.organization),
|
||||
)
|
||||
Message.objects.create(
|
||||
conversation=conversation, author_type=MessageAuthor.AI, text=fallback
|
||||
)
|
||||
_finish(turn.message, AiTurnState.FAILED)
|
||||
|
||||
notify(
|
||||
context=context,
|
||||
type=NotificationType.OPERATOR_REQUESTED,
|
||||
audience=NotificationAudience.OPERATORS,
|
||||
audience_group=conversation.group,
|
||||
title=f"Нужен оператор · {_contact_name(turn)}",
|
||||
title_key="notifications.operator_needed",
|
||||
text_params={"contact": _contact_name(turn)},
|
||||
body="AI временно недоступен, диалог ждёт ответа",
|
||||
body_key="notifications.ai_unavailable_waiting",
|
||||
target_id=conversation.id,
|
||||
source_type="Conversation",
|
||||
source_id=conversation.id,
|
||||
dedup_key=f"aifail:{conversation.id}",
|
||||
)
|
||||
notify_management(
|
||||
context=context,
|
||||
type=NotificationType.AI_STOPPED,
|
||||
title=f"Ошибка AI · {channel.name}",
|
||||
body="AI временно недоступен, диалог передан оператору",
|
||||
title_key="notifications.ai_error",
|
||||
body_key="notifications.ai_unavailable_handed_over",
|
||||
text_params={"channel": channel.name},
|
||||
target_id=conversation.id,
|
||||
source_type="Conversation",
|
||||
source_id=conversation.id,
|
||||
dedup_key=f"aierror:{conversation.id}",
|
||||
)
|
||||
return fallback
|
||||
@@ -10,3 +10,4 @@ class ConversationsConfig(AppConfig):
|
||||
def ready(self) -> None:
|
||||
# Свежесть диалога поддерживает сигнал: сообщения создаются в семи местах.
|
||||
from chatballs.conversations import signals # noqa: F401
|
||||
from chatballs.conversations import event_handlers # noqa: F401 (register outbox handlers)
|
||||
@@ -0,0 +1,14 @@
|
||||
"""Обработчики outbox-событий домена диалогов."""
|
||||
|
||||
from chatballs.conversations.ai_turn import AI_TURN_REQUESTED, run_requested_turn
|
||||
from chatballs.events.handlers import register
|
||||
from chatballs.tenancy.context import TenantContext
|
||||
|
||||
|
||||
@register(AI_TURN_REQUESTED, manages_own_transaction=True)
|
||||
def handle_ai_turn_requested(payload: dict, context: TenantContext | None) -> None:
|
||||
"""Ход AI сам управляет транзакциями: он ходит к провайдеру и в мессенджер,
|
||||
и держать ради этого одну транзакцию на весь обработчик нельзя."""
|
||||
if context is None: # pragma: no cover - событие диалога всегда арендное
|
||||
return
|
||||
run_requested_turn(payload, context)
|
||||
@@ -1,23 +1,23 @@
|
||||
"""Inbound ingest for messenger connections (M2a).
|
||||
"""Приём входящих из подключений (M2a).
|
||||
|
||||
One inbound message -> contact/conversation/message -> AI turn (if the dialog is
|
||||
AI-controlled) -> outbound reply. Idempotent via the events InboxEvent.
|
||||
Одно входящее -> контакт/диалог/сообщение -> заявка на ход AI, если диалог
|
||||
ведёт агент. Повторы отсекаются через InboxEvent.
|
||||
|
||||
Обращений наружу здесь нет и быть не должно: приём вызывают цикл опроса
|
||||
мессенджеров и HTTP-запрос виджета, и ждать провайдера ни тот, ни другой не
|
||||
может. Ответ считает роль событий (chatballs.conversations.ai_turn).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import logging
|
||||
from dataclasses import dataclass
|
||||
|
||||
from django.db import IntegrityError, transaction
|
||||
from django.utils import timezone
|
||||
|
||||
from chatballs.ai.models import HISTORY_LIMIT_DEFAULT
|
||||
from chatballs.ai.provider.base import ProviderError
|
||||
from chatballs.ai.runtime import HANDOFF_TOKEN
|
||||
from chatballs.channels.runtime import run_channel_turn
|
||||
from chatballs.conversations import transports
|
||||
from chatballs.conversations.ai_turn import request_ai_turn
|
||||
from chatballs.conversations.contact_avatars import refresh_contact_avatar
|
||||
from chatballs.conversations.models import (
|
||||
ConnectionIdentity,
|
||||
@@ -29,26 +29,17 @@ from chatballs.conversations.models import (
|
||||
Message,
|
||||
MessageAuthor,
|
||||
MessageKind,
|
||||
SystemEvent,
|
||||
TranscriptStatus,
|
||||
)
|
||||
from chatballs.conversations.queue import QUEUE_FIELDS, enter_queue, is_waiting
|
||||
from chatballs.conversations.transports.base import InboundMessage
|
||||
from chatballs.events.models import EventOwnership, InboxEvent
|
||||
from chatballs.i18n import t
|
||||
from chatballs.notifications.models import NotificationAudience, NotificationType
|
||||
from chatballs.notifications.services import notify, notify_management
|
||||
from chatballs.notifications.services import notify
|
||||
from chatballs.tenancy.context import TenantContext
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_ROLE = {
|
||||
MessageAuthor.CONTACT: "user",
|
||||
MessageAuthor.AI: "assistant",
|
||||
MessageAuthor.OPERATOR: "assistant",
|
||||
MessageAuthor.SYSTEM: "system",
|
||||
}
|
||||
|
||||
|
||||
def _already_processed(context: TenantContext, source: str, external_id: str, text: str) -> bool:
|
||||
"""Отметить сообщение обработанным; True — оно уже приходило.
|
||||
@@ -74,105 +65,6 @@ def _already_processed(context: TenantContext, source: str, external_id: str, te
|
||||
return True
|
||||
|
||||
|
||||
def _history(conversation: Conversation, limit: int) -> list[dict]:
|
||||
# С конца и с ограничением в базе: длинный диалог не поднимается в память
|
||||
# целиком ради последних сообщений. Самое новое — только что сохранённое
|
||||
# входящее, оно уходит модели отдельно.
|
||||
latest = conversation.messages.order_by("-created_at", "-id")[: limit + 1]
|
||||
prior = list(reversed(latest))[:-1]
|
||||
# Голосовые попадают в контекст стенограммой.
|
||||
return [{"role": _ROLE.get(m.author_type, "user"), "content": m.text or m.transcript} for m in prior if m.text or m.transcript]
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class TranscriptionJob:
|
||||
"""Всё, что нужно провайдеру, — уже прочитанное из базы и хранилища.
|
||||
|
||||
Разложено на три шага (``prepare`` → ``run`` → ``store``), чтобы вызывающий
|
||||
мог держать транзакцию только вокруг первого и третьего: обращение к
|
||||
провайдеру ждёт ответа десятки секунд, и всё это время транзакция занимала
|
||||
бы соединение из пула (chatballs.tenancy.middleware).
|
||||
"""
|
||||
|
||||
provider: object
|
||||
model: str
|
||||
audio: bytes
|
||||
filename: str
|
||||
content_type: str
|
||||
|
||||
|
||||
def prepare_transcription(channel, message: Message) -> TranscriptionJob | None:
|
||||
"""Шаг в транзакции: провайдер организации, модель и байты аудио."""
|
||||
from chatballs.ai.provider.factory import get_transcription_provider
|
||||
from chatballs.ai.provider.routing import (
|
||||
DEFAULT_TRANSCRIPTION_MODEL,
|
||||
resolve_transcription_model,
|
||||
)
|
||||
|
||||
if not message.audio:
|
||||
return None
|
||||
provider = get_transcription_provider(channel=channel)
|
||||
try:
|
||||
model = resolve_transcription_model(channel)
|
||||
except ProviderError:
|
||||
model = DEFAULT_TRANSCRIPTION_MODEL # тестовый провайдер без интеграции
|
||||
with message.audio.open("rb") as handle:
|
||||
audio = handle.read()
|
||||
return TranscriptionJob(
|
||||
provider=provider,
|
||||
model=model,
|
||||
audio=audio,
|
||||
filename=message.audio.name.rsplit("/", 1)[-1],
|
||||
content_type=message.audio_content_type or "audio/ogg",
|
||||
)
|
||||
|
||||
|
||||
def run_transcription(job: TranscriptionJob) -> str:
|
||||
"""Шаг без транзакции: обращение к провайдеру."""
|
||||
return job.provider.transcribe(
|
||||
audio=job.audio,
|
||||
filename=job.filename,
|
||||
content_type=job.content_type,
|
||||
model=job.model,
|
||||
).strip()
|
||||
|
||||
|
||||
def store_transcription(message: Message, transcript: str) -> None:
|
||||
"""Шаг в транзакции: сохранить стенограмму и статус."""
|
||||
message.transcript = transcript
|
||||
message.transcript_status = TranscriptStatus.READY if transcript else TranscriptStatus.FAILED
|
||||
message.save(update_fields=["transcript", "transcript_status"])
|
||||
|
||||
|
||||
def mark_transcription_failed(message: Message) -> None:
|
||||
"""Статус FAILED — оператор повторит кнопкой."""
|
||||
message.transcript_status = TranscriptStatus.FAILED
|
||||
message.save(update_fields=["transcript_status"])
|
||||
|
||||
|
||||
def transcribe_voice_message(channel, message: Message, *, raise_errors: bool = False) -> str:
|
||||
"""Стенограмма голосового через BYOK-провайдера организации; пустая строка,
|
||||
если провайдер не умеет или недоступен (статус FAILED — оператор повторит кнопкой).
|
||||
|
||||
Три шага подряд, в транзакции вызывающего: так входящее сообщение
|
||||
обрабатывается целиком (ingest_inbound). Оператору, нажавшему «расшифровать»,
|
||||
ждать под транзакцией незачем — там шаги разнесены (voice_views).
|
||||
"""
|
||||
try:
|
||||
job = prepare_transcription(channel, message)
|
||||
if job is None:
|
||||
return ""
|
||||
transcript = run_transcription(job)
|
||||
except ProviderError as error:
|
||||
logger.info("Voice transcription unavailable for message %s: %s", message.id, error)
|
||||
mark_transcription_failed(message)
|
||||
if raise_errors:
|
||||
raise
|
||||
return ""
|
||||
store_transcription(message, transcript)
|
||||
return transcript
|
||||
|
||||
|
||||
def ingest_inbound(integration, inbound: InboundMessage) -> None:
|
||||
channel = integration.channel
|
||||
if channel is None:
|
||||
@@ -374,13 +266,10 @@ def ingest_inbound(integration, inbound: InboundMessage) -> None:
|
||||
transports.send_contact_ack(integration, chat_id=conversation.external_chat_id, user_id=inbound.user_id, text=ack)
|
||||
return
|
||||
|
||||
# Голосовое: AI отвечает текстом по стенограмме (BYOK-провайдер). Если
|
||||
# расшифровка недоступна, а также для файлов без текста — диалог уходит
|
||||
# оператору, как при недоступном AI, но без имитации сбоя.
|
||||
ai_input = inbound.text
|
||||
if is_voice and conversation.control_mode == ControlMode.AI and ai_available:
|
||||
ai_input = transcribe_voice_message(channel, message)
|
||||
if (is_voice and not ai_input) or files_only:
|
||||
# Файлы без текста: отвечать не на что — диалог уходит оператору, как при
|
||||
# недоступном AI, но без имитации сбоя. Голосовое сюда не попадает: его
|
||||
# расшифровка — это обращение к провайдеру, и она идёт ходом AI.
|
||||
if files_only:
|
||||
if conversation.control_mode == ControlMode.AI:
|
||||
enter_queue(conversation)
|
||||
conversation.save(update_fields=QUEUE_FIELDS)
|
||||
@@ -394,8 +283,7 @@ def ingest_inbound(integration, inbound: InboundMessage) -> None:
|
||||
title=f"Нужен оператор · {contact.name or 'Гость'}",
|
||||
title_key="notifications.operator_needed",
|
||||
text_params={"contact": contact.name or t("conversations.guest")},
|
||||
body="Голосовое без расшифровки" if is_voice else message_text[:120],
|
||||
body_key="notifications.voice_without_transcript" if is_voice else "",
|
||||
body=message_text[:120],
|
||||
target_id=conversation.id,
|
||||
source_type="Conversation",
|
||||
source_id=conversation.id,
|
||||
@@ -408,99 +296,15 @@ def ingest_inbound(integration, inbound: InboundMessage) -> None:
|
||||
if conversation.control_mode != ControlMode.AI:
|
||||
return
|
||||
|
||||
try:
|
||||
result = run_channel_turn(
|
||||
channel=channel,
|
||||
message=ai_input,
|
||||
history=_history(
|
||||
conversation, agent.history_limit if agent else HISTORY_LIMIT_DEFAULT
|
||||
),
|
||||
)
|
||||
except ProviderError as error:
|
||||
# Сбой AI не должен «терять» сообщение: переводим диалог в очередь к
|
||||
# оператору, уведомляем и отвечаем клиенту понятным fallback.
|
||||
logger.warning("AI turn failed for conversation %s: %s", conversation.id, error)
|
||||
enter_queue(conversation)
|
||||
conversation.last_activity_at = timezone.now()
|
||||
conversation.save(update_fields=[*QUEUE_FIELDS, "last_activity_at"])
|
||||
Message.objects.create(
|
||||
conversation=conversation,
|
||||
author_type=MessageAuthor.SYSTEM,
|
||||
system_event=SystemEvent.AI_UNAVAILABLE,
|
||||
text="AI недоступен — диалог передан оператору",
|
||||
)
|
||||
fallback = "Извините, прямо сейчас не получается ответить. Я передал ваш вопрос специалисту — он скоро подключится."
|
||||
Message.objects.create(conversation=conversation, author_type=MessageAuthor.AI, text=fallback)
|
||||
notify(
|
||||
context=context,
|
||||
type=NotificationType.OPERATOR_REQUESTED,
|
||||
audience=NotificationAudience.OPERATORS,
|
||||
audience_group=conversation.group,
|
||||
title=f"Нужен оператор · {contact.name or 'Гость'}",
|
||||
title_key="notifications.operator_needed",
|
||||
text_params={"contact": contact.name or t("conversations.guest")},
|
||||
body="AI временно недоступен, диалог ждёт ответа",
|
||||
body_key="notifications.ai_unavailable_waiting",
|
||||
target_id=conversation.id,
|
||||
source_type="Conversation",
|
||||
source_id=conversation.id,
|
||||
dedup_key=f"aifail:{conversation.id}",
|
||||
)
|
||||
notify_management(
|
||||
context=context,
|
||||
type=NotificationType.AI_STOPPED,
|
||||
title=f"Ошибка AI · {channel.name}",
|
||||
body="AI временно недоступен, диалог передан оператору",
|
||||
title_key="notifications.ai_error",
|
||||
body_key="notifications.ai_unavailable_handed_over",
|
||||
text_params={"channel": channel.name},
|
||||
target_id=conversation.id,
|
||||
source_type="Conversation",
|
||||
source_id=conversation.id,
|
||||
dedup_key=f"aierror:{conversation.id}",
|
||||
)
|
||||
transports.send_reply(integration, chat_id=conversation.external_chat_id, user_id=inbound.user_id, text=fallback)
|
||||
return
|
||||
|
||||
reply = result.text
|
||||
handoff = HANDOFF_TOKEN in reply
|
||||
if handoff:
|
||||
reply = reply.replace(HANDOFF_TOKEN, "").strip()
|
||||
|
||||
Message.objects.create(conversation=conversation, author_type=MessageAuthor.AI, text=reply)
|
||||
conversation.last_activity_at = timezone.now()
|
||||
if handoff:
|
||||
enter_queue(conversation)
|
||||
else:
|
||||
conversation.expected_responder = ExpectedResponder.CUSTOMER
|
||||
conversation.save(update_fields=[*QUEUE_FIELDS, "last_activity_at"])
|
||||
|
||||
if handoff:
|
||||
Message.objects.create(
|
||||
conversation=conversation,
|
||||
author_type=MessageAuthor.SYSTEM,
|
||||
system_event=SystemEvent.AI_HANDED_OVER,
|
||||
text="AI передал диалог оператору",
|
||||
)
|
||||
notify(
|
||||
context=context,
|
||||
type=NotificationType.OPERATOR_REQUESTED,
|
||||
audience=NotificationAudience.OPERATORS,
|
||||
audience_group=conversation.group,
|
||||
title=f"AI передал диалог · {contact.name or 'Гость'}",
|
||||
title_key="notifications.ai_handed_over",
|
||||
text_params={"contact": contact.name or t("conversations.guest")},
|
||||
body=ai_input[:120],
|
||||
target_id=conversation.id,
|
||||
source_type="Conversation",
|
||||
source_id=conversation.id,
|
||||
dedup_key=f"handoff:{conversation.id}",
|
||||
)
|
||||
|
||||
if reply:
|
||||
transports.send_reply(
|
||||
integration, chat_id=conversation.external_chat_id, user_id=inbound.user_id, text=reply
|
||||
)
|
||||
# Ход AI — отдельная работа: обращение к модели ждёт ответа секунды и
|
||||
# десятки секунд, а приём входящих столько ждать не может. Здесь только
|
||||
# заявка; считает ход роль событий (chatballs.conversations.ai_turn).
|
||||
request_ai_turn(
|
||||
message=message,
|
||||
user_id=inbound.user_id,
|
||||
context=context,
|
||||
is_new_conversation=is_new,
|
||||
)
|
||||
|
||||
|
||||
def _store_attachment(integration, inbound_file, message: Message) -> None:
|
||||
|
||||
@@ -0,0 +1,18 @@
|
||||
# Generated by Django 5.2.16 on 2026-09-20 02:00
|
||||
|
||||
from django.db import migrations, models
|
||||
|
||||
|
||||
class Migration(migrations.Migration):
|
||||
|
||||
dependencies = [
|
||||
('conversations', '0025_contact_avatar_contact_avatar_content_type_and_more'),
|
||||
]
|
||||
|
||||
operations = [
|
||||
migrations.AddField(
|
||||
model_name='message',
|
||||
name='ai_turn_state',
|
||||
field=models.CharField(choices=[('NONE', 'Ход не нужен'), ('PENDING', 'Ожидает'), ('RUNNING', 'Считается'), ('DONE', 'Отвечено'), ('FAILED', 'Не удалось')], default='NONE', max_length=8),
|
||||
),
|
||||
]
|
||||
@@ -324,6 +324,22 @@ class MessageKind(models.TextChoices):
|
||||
FILE = "file", "Файл"
|
||||
|
||||
|
||||
class AiTurnState(models.TextChoices):
|
||||
"""Состояние хода AI по входящему сообщению.
|
||||
|
||||
Ответ считается не в приёме, а отдельной ролью воркера
|
||||
(chatballs.conversations.ai_turn), поэтому у входящего появилось состояние.
|
||||
По нему видно, что ответ ещё считается — виджет показывает «печатает», — и
|
||||
по нему же повторная доставка события не приводит ко второму ответу.
|
||||
"""
|
||||
|
||||
NONE = "NONE", "Ход не нужен"
|
||||
PENDING = "PENDING", "Ожидает"
|
||||
RUNNING = "RUNNING", "Считается"
|
||||
DONE = "DONE", "Отвечено"
|
||||
FAILED = "FAILED", "Не удалось"
|
||||
|
||||
|
||||
class TranscriptStatus(models.TextChoices):
|
||||
# Расшифровка голосового (дизайн-базлайн v2, кадр H): по кнопке, через
|
||||
# BYOK-провайдера организации (решение владельца 2026-09-04).
|
||||
@@ -375,6 +391,10 @@ class Message(TenantRelationModel):
|
||||
transcript_status = models.CharField(
|
||||
max_length=8, choices=TranscriptStatus.choices, default=TranscriptStatus.NONE
|
||||
)
|
||||
# Ход AI по этому сообщению: ожидает, считается, отвечено, не удалось.
|
||||
ai_turn_state = models.CharField(
|
||||
max_length=8, choices=AiTurnState.choices, default=AiTurnState.NONE
|
||||
)
|
||||
# Файл/фото (kind=FILE): вложение с исходным именем, типом и размером.
|
||||
attachment = models.FileField(upload_to=message_attachment_upload_path, max_length=512, blank=True)
|
||||
attachment_name = models.CharField(max_length=255, blank=True)
|
||||
|
||||
@@ -0,0 +1,141 @@
|
||||
"""Ход AI как отдельная работа: приём не ждёт модель, ответ считается событием."""
|
||||
|
||||
from unittest import mock
|
||||
|
||||
from django.test import TestCase, override_settings
|
||||
|
||||
from chatballs.ai.models import AIAgent, AIAgentStatus
|
||||
from chatballs.channels.models import Channel
|
||||
from chatballs.conversations.ai_turn import AI_TURN_REQUESTED
|
||||
from chatballs.conversations.ingest import ingest_inbound
|
||||
from chatballs.conversations.models import (
|
||||
AiTurnState,
|
||||
ControlMode,
|
||||
ExpectedResponder,
|
||||
Message,
|
||||
MessageAuthor,
|
||||
)
|
||||
from chatballs.conversations.transports.base import InboundMessage
|
||||
from chatballs.events.handlers import dispatch
|
||||
from chatballs.events.models import OutboxEvent
|
||||
from chatballs.identity.bootstrap import bootstrap_owner
|
||||
from chatballs.identity.models import Organization
|
||||
from chatballs.integrations.models import Integration, IntegrationKind, IntegrationProvider
|
||||
from chatballs.testing import ai_answer, run_pending_ai_turns
|
||||
|
||||
|
||||
class AiTurnQueueTests(TestCase):
|
||||
def setUp(self) -> None:
|
||||
bootstrap_owner(email="owner@example.com", password="temporary-password")
|
||||
self.organization = Organization.objects.get(slug="demo")
|
||||
self.channel = Channel.objects.create(
|
||||
organization=self.organization, code="line", name="Линия"
|
||||
)
|
||||
AIAgent.objects.create(
|
||||
channel=self.channel,
|
||||
name="Агент",
|
||||
model="openai/gpt-4o-mini",
|
||||
status=AIAgentStatus.ACTIVE,
|
||||
)
|
||||
self.integration = Integration.objects.create(
|
||||
organization=self.organization,
|
||||
kind=IntegrationKind.MESSENGER,
|
||||
provider=IntegrationProvider.TELEGRAM,
|
||||
name="Bot",
|
||||
secret="token",
|
||||
channel=self.channel,
|
||||
)
|
||||
self.inbound = InboundMessage(
|
||||
external_id="ext-1",
|
||||
user_id="u-1",
|
||||
chat_id="c-1",
|
||||
text="Здравствуйте",
|
||||
display_name="Гость",
|
||||
)
|
||||
|
||||
def _inbound_message(self) -> Message:
|
||||
return Message.objects.get(author_type=MessageAuthor.CONTACT)
|
||||
|
||||
def _ai_messages(self):
|
||||
return Message.objects.filter(author_type=MessageAuthor.AI)
|
||||
|
||||
def test_ingest_queues_the_turn_and_does_not_call_the_model(self) -> None:
|
||||
# Главное свойство всей развязки: приём не ждёт провайдера.
|
||||
with mock.patch("chatballs.ai.provider.local.LocalProvider.chat") as chat:
|
||||
ingest_inbound(self.integration, self.inbound)
|
||||
|
||||
chat.assert_not_called()
|
||||
message = self._inbound_message()
|
||||
self.assertEqual(message.ai_turn_state, AiTurnState.PENDING)
|
||||
self.assertFalse(self._ai_messages().exists())
|
||||
self.assertTrue(
|
||||
OutboxEvent.objects.filter(
|
||||
event_type=AI_TURN_REQUESTED, aggregate_id=str(message.conversation_id)
|
||||
).exists()
|
||||
)
|
||||
|
||||
def test_turn_answers_and_closes_the_message(self) -> None:
|
||||
with (
|
||||
ai_answer("Здравствуйте!"),
|
||||
mock.patch(
|
||||
"chatballs.conversations.transports.send_reply", return_value=True
|
||||
) as send,
|
||||
):
|
||||
ingest_inbound(self.integration, self.inbound)
|
||||
self.assertEqual(run_pending_ai_turns(), 1)
|
||||
|
||||
self.assertEqual(self._ai_messages().get().text, "Здравствуйте!")
|
||||
self.assertEqual(self._inbound_message().ai_turn_state, AiTurnState.DONE)
|
||||
conversation = self.channel.conversations.get()
|
||||
self.assertEqual(conversation.control_mode, ControlMode.AI)
|
||||
self.assertEqual(conversation.expected_responder, ExpectedResponder.CUSTOMER)
|
||||
send.assert_called_once()
|
||||
|
||||
def test_repeated_delivery_does_not_answer_twice(self) -> None:
|
||||
# Событие могут привезти второй раз: процесс упал между ответом и
|
||||
# отметкой о нём. Второй ответ клиенту — это хуже, чем ни одного.
|
||||
with (
|
||||
ai_answer("Здравствуйте!"),
|
||||
mock.patch("chatballs.conversations.transports.send_reply", return_value=True),
|
||||
):
|
||||
ingest_inbound(self.integration, self.inbound)
|
||||
event = OutboxEvent.objects.get(event_type=AI_TURN_REQUESTED)
|
||||
dispatch(event)
|
||||
dispatch(event)
|
||||
|
||||
self.assertEqual(self._ai_messages().count(), 1)
|
||||
|
||||
def test_turn_for_a_dialog_taken_by_an_operator_is_dropped(self) -> None:
|
||||
with (
|
||||
ai_answer("Здравствуйте!"),
|
||||
mock.patch("chatballs.conversations.transports.send_reply", return_value=True),
|
||||
):
|
||||
ingest_inbound(self.integration, self.inbound)
|
||||
conversation = self.channel.conversations.get()
|
||||
conversation.control_mode = ControlMode.HUMAN
|
||||
conversation.save(update_fields=["control_mode"])
|
||||
run_pending_ai_turns()
|
||||
|
||||
self.assertFalse(self._ai_messages().exists())
|
||||
self.assertEqual(self._inbound_message().ai_turn_state, AiTurnState.DONE)
|
||||
|
||||
@override_settings(CHATBALLS_AI_TURN_DEADLINE_SECONDS=0)
|
||||
def test_expired_turn_goes_to_the_operator_instead_of_the_model(self) -> None:
|
||||
# Ответ, пролежавший в очереди, клиенту уже не нужен — нужен человек.
|
||||
with (
|
||||
mock.patch("chatballs.ai.provider.local.LocalProvider.chat") as chat,
|
||||
mock.patch(
|
||||
"chatballs.conversations.transports.send_reply", return_value=True
|
||||
) as send,
|
||||
):
|
||||
ingest_inbound(self.integration, self.inbound)
|
||||
run_pending_ai_turns()
|
||||
|
||||
chat.assert_not_called()
|
||||
conversation = self.channel.conversations.get()
|
||||
self.assertEqual(conversation.control_mode, ControlMode.PAUSED)
|
||||
self.assertEqual(conversation.expected_responder, ExpectedResponder.OPERATOR)
|
||||
self.assertEqual(self._inbound_message().ai_turn_state, AiTurnState.FAILED)
|
||||
self.assertTrue(self._ai_messages().filter(text__contains="специалисту").exists())
|
||||
# Клиент получает этот текст в своём канале, а не только в базе.
|
||||
send.assert_called_once()
|
||||
@@ -20,6 +20,8 @@ from chatballs.identity.bootstrap import bootstrap_owner
|
||||
from chatballs.identity.models import Organization
|
||||
from chatballs.integrations.models import Integration, IntegrationKind, IntegrationProvider
|
||||
|
||||
from chatballs.testing import ai_answer, run_pending_ai_turns
|
||||
|
||||
EMAIL_CONFIG = {
|
||||
|
||||
"email": "support@example.com",
|
||||
@@ -505,14 +507,16 @@ class EmailIngestThreadMetaTests(TestCase):
|
||||
|
||||
with (
|
||||
|
||||
mock.patch("chatballs.conversations.ingest.run_channel_turn", return_value=mock.Mock(text="Ответ")),
|
||||
ai_answer("Ответ"),
|
||||
|
||||
mock.patch("chatballs.conversations.ingest.transports.send_reply", return_value=True),
|
||||
mock.patch("chatballs.conversations.transports.send_reply", return_value=True),
|
||||
|
||||
):
|
||||
|
||||
ingest_inbound(self.integration, inbound)
|
||||
|
||||
run_pending_ai_turns()
|
||||
|
||||
|
||||
|
||||
def test_subject_pinned_to_first_message_id_follows_last(self) -> None:
|
||||
|
||||
@@ -22,7 +22,7 @@ from chatballs.integrations.models import (
|
||||
IntegrationKind,
|
||||
IntegrationProvider,
|
||||
)
|
||||
from chatballs.testing import TenantAPIClient as APIClient
|
||||
from chatballs.testing import TenantAPIClient as APIClient, run_pending_ai_turns
|
||||
|
||||
|
||||
def _connection(channel: Channel) -> Integration:
|
||||
@@ -71,33 +71,21 @@ class OperatorOnlyIngestTests(TestCase):
|
||||
|
||||
|
||||
|
||||
def _ingest_without_ai(self, inbound: InboundMessage) -> tuple[mock.Mock, mock.Mock]:
|
||||
def _ingest_without_ai(self, inbound: InboundMessage) -> tuple[int, mock.Mock]:
|
||||
|
||||
with (
|
||||
|
||||
mock.patch(
|
||||
|
||||
"chatballs.conversations.ingest.run_channel_turn"
|
||||
|
||||
) as ai_turn,
|
||||
|
||||
mock.patch(
|
||||
|
||||
"chatballs.conversations.ingest.transports.send_reply"
|
||||
|
||||
) as send,
|
||||
|
||||
):
|
||||
with mock.patch("chatballs.conversations.transports.send_reply") as send:
|
||||
|
||||
ingest_inbound(self.integration, inbound)
|
||||
|
||||
return ai_turn, send
|
||||
turns = run_pending_ai_turns()
|
||||
|
||||
return turns, send
|
||||
|
||||
|
||||
|
||||
def test_new_dialog_starts_in_queue_without_ai_fallback(self) -> None:
|
||||
|
||||
ai_turn, send = self._ingest_without_ai(
|
||||
turns, send = self._ingest_without_ai(
|
||||
|
||||
InboundMessage(
|
||||
|
||||
@@ -135,7 +123,7 @@ class OperatorOnlyIngestTests(TestCase):
|
||||
|
||||
)
|
||||
|
||||
ai_turn.assert_not_called()
|
||||
self.assertEqual(turns, 0)
|
||||
|
||||
send.assert_not_called()
|
||||
|
||||
@@ -179,7 +167,7 @@ class OperatorOnlyIngestTests(TestCase):
|
||||
|
||||
|
||||
|
||||
ai_turn, send = self._ingest_without_ai(
|
||||
turns, send = self._ingest_without_ai(
|
||||
|
||||
InboundMessage(
|
||||
|
||||
@@ -211,7 +199,7 @@ class OperatorOnlyIngestTests(TestCase):
|
||||
|
||||
self.assertEqual(conversation.messages.count(), 1)
|
||||
|
||||
ai_turn.assert_not_called()
|
||||
self.assertEqual(turns, 0)
|
||||
|
||||
send.assert_not_called()
|
||||
|
||||
|
||||
@@ -26,7 +26,7 @@ from chatballs.identity.bootstrap import bootstrap_owner
|
||||
from chatballs.identity.models import HumanUser, Organization
|
||||
from chatballs.integrations.models import Integration, IntegrationKind, IntegrationProvider
|
||||
from chatballs.notifications.models import Notification, NotificationType
|
||||
from chatballs.testing import tenant_context_for
|
||||
from chatballs.testing import ai_answer, tenant_context_for
|
||||
|
||||
|
||||
class QueueTestBase(TestCase):
|
||||
@@ -55,10 +55,7 @@ class QueueTestBase(TestCase):
|
||||
text=text,
|
||||
display_name=chat_id,
|
||||
)
|
||||
with (
|
||||
mock.patch("chatballs.conversations.ingest.run_channel_turn"),
|
||||
mock.patch("chatballs.conversations.ingest.transports.send_reply"),
|
||||
):
|
||||
with mock.patch("chatballs.conversations.transports.send_reply"):
|
||||
ingest_inbound(self.integration, inbound)
|
||||
|
||||
def _waiting_order(self) -> list[int]:
|
||||
@@ -155,10 +152,9 @@ class NewDialogNotificationTests(QueueTestBase):
|
||||
is_active=True,
|
||||
)
|
||||
self.assertTrue(agent.is_active)
|
||||
with mock.patch(
|
||||
"chatballs.conversations.ingest.run_channel_turn",
|
||||
return_value=mock.Mock(text="Здравствуйте!"),
|
||||
), mock.patch("chatballs.conversations.ingest.transports.send_reply"):
|
||||
with ai_answer("Здравствуйте!"), mock.patch(
|
||||
"chatballs.conversations.transports.send_reply"
|
||||
):
|
||||
ingest_inbound(
|
||||
self.integration,
|
||||
InboundMessage(
|
||||
|
||||
@@ -16,7 +16,8 @@ from chatballs.identity.bootstrap import bootstrap_owner
|
||||
from chatballs.identity.models import Organization
|
||||
from chatballs.integrations.models import Integration, IntegrationKind, IntegrationProvider
|
||||
from chatballs.tenancy.database import tenant_atomic
|
||||
from chatballs.testing import TenantAPIClient as APIClient
|
||||
from chatballs.conversations import ai_turn
|
||||
from chatballs.testing import TenantAPIClient as APIClient, ai_answer, run_pending_ai_turns
|
||||
|
||||
|
||||
class VoiceAiReplyTests(TestCase):
|
||||
@@ -30,43 +31,53 @@ class VoiceAiReplyTests(TestCase):
|
||||
)
|
||||
self.inbound = InboundMessage(external_id="v-1", user_id="u-1", chat_id="c-1", text="", display_name="Ольга", voice_file_id="f-1", voice_duration=5, voice_mime="audio/ogg")
|
||||
|
||||
def _ingest(self, transcribe, turn):
|
||||
def _ingest(self, transcribe, answer="Ответ"):
|
||||
"""Приём голосового и ход AI по нему.
|
||||
|
||||
Расшифровка — обращение к провайдеру, поэтому она идёт не в приёме, а
|
||||
в ходе (chatballs.conversations.ai_turn); тест повторяет этот порядок.
|
||||
"""
|
||||
with (
|
||||
mock.patch("chatballs.conversations.ingest.transports.download_voice", return_value=(b"OGG", "audio/ogg")),
|
||||
mock.patch("chatballs.conversations.ingest.transports.send_reply", return_value=True) as send,
|
||||
mock.patch("chatballs.conversations.transports.send_reply", return_value=True) as send,
|
||||
mock.patch("chatballs.ai.provider.local.LocalProvider.transcribe", **transcribe),
|
||||
mock.patch("chatballs.conversations.ingest.run_channel_turn", **turn) as run,
|
||||
mock.patch("chatballs.conversations.ai_turn.plan_chat", wraps=ai_turn.plan_chat) as plan,
|
||||
ai_answer(answer),
|
||||
tenant_atomic(self.organization.id),
|
||||
):
|
||||
ingest_inbound(self.integration, self.inbound)
|
||||
return send, run
|
||||
run_pending_ai_turns()
|
||||
return send, plan
|
||||
|
||||
def test_ai_answers_voice_by_transcript(self) -> None:
|
||||
send, run = self._ingest({"return_value": "Можно оформить возврат?"}, {"return_value": mock.Mock(text="Да, возврат возможен в течение 14 дней.")})
|
||||
send, plan = self._ingest(
|
||||
{"return_value": "Можно оформить возврат?"},
|
||||
answer="Да, возврат возможен в течение 14 дней.",
|
||||
)
|
||||
conversation = self.channel.conversations.get()
|
||||
voice = conversation.messages.get(kind=MessageKind.VOICE)
|
||||
self.assertEqual(voice.transcript, "Можно оформить возврат?")
|
||||
self.assertEqual(voice.transcript_status, TranscriptStatus.READY)
|
||||
run.assert_called_once()
|
||||
self.assertEqual(run.call_args.kwargs["message"], "Можно оформить возврат?")
|
||||
plan.assert_called_once()
|
||||
self.assertEqual(plan.call_args.kwargs["message"], "Можно оформить возврат?")
|
||||
reply = conversation.messages.get(author_type=MessageAuthor.AI)
|
||||
self.assertIn("возврат", reply.text)
|
||||
send.assert_called_once()
|
||||
self.assertEqual(conversation.control_mode, ControlMode.AI)
|
||||
|
||||
def test_without_transcription_dialog_goes_to_operator(self) -> None:
|
||||
send, run = self._ingest({"side_effect": ProviderError("нет STT")}, {"return_value": mock.Mock(text="x")})
|
||||
send, plan = self._ingest({"side_effect": ProviderError("нет STT")})
|
||||
conversation = self.channel.conversations.get()
|
||||
voice = conversation.messages.get(kind=MessageKind.VOICE)
|
||||
self.assertEqual(voice.transcript_status, TranscriptStatus.FAILED)
|
||||
run.assert_not_called()
|
||||
plan.assert_not_called()
|
||||
send.assert_not_called()
|
||||
self.assertEqual(conversation.control_mode, ControlMode.PAUSED)
|
||||
|
||||
def test_transcript_is_in_ai_history(self) -> None:
|
||||
from chatballs.conversations.ingest import _history
|
||||
from chatballs.conversations.ai_turn import _history
|
||||
|
||||
self._ingest({"return_value": "Первый вопрос"}, {"return_value": mock.Mock(text="Ответ")})
|
||||
self._ingest({"return_value": "Первый вопрос"})
|
||||
conversation = self.channel.conversations.get()
|
||||
conversation.messages.create(author_type=MessageAuthor.CONTACT, text="Второй")
|
||||
roles = [(h["role"], h["content"]) for h in _history(conversation, 20)]
|
||||
@@ -75,16 +86,16 @@ class VoiceAiReplyTests(TestCase):
|
||||
def test_ai_history_window_follows_agent_setting(self) -> None:
|
||||
self.channel.ai_agent.history_limit = 3
|
||||
self.channel.ai_agent.save(update_fields=["history_limit"])
|
||||
self._ingest({"return_value": "Первый вопрос"}, {"return_value": mock.Mock(text="Ответ")})
|
||||
self._ingest({"return_value": "Первый вопрос"})
|
||||
conversation = self.channel.conversations.get()
|
||||
for number in range(1, 6):
|
||||
conversation.messages.create(author_type=MessageAuthor.CONTACT, text=f"Сообщение {number}")
|
||||
self.inbound = InboundMessage(external_id="t-2", user_id="u-1", chat_id="c-1", text="Последнее", display_name="Ольга")
|
||||
_, run = self._ingest({"return_value": ""}, {"return_value": mock.Mock(text="Ответ")})
|
||||
history = [item["content"] for item in run.call_args.kwargs["history"]]
|
||||
_, plan = self._ingest({"return_value": ""})
|
||||
history = [item["content"] for item in plan.call_args.kwargs["history"]]
|
||||
# Три сообщения перед новым; само новое уходит модели отдельно.
|
||||
self.assertEqual(history, ["Сообщение 3", "Сообщение 4", "Сообщение 5"])
|
||||
self.assertEqual(run.call_args.kwargs["message"], "Последнее")
|
||||
self.assertEqual(plan.call_args.kwargs["message"], "Последнее")
|
||||
|
||||
|
||||
class CommunicationSettingsTests(TestCase):
|
||||
|
||||
@@ -4,7 +4,6 @@ from unittest import mock
|
||||
from django.test import TestCase, override_settings
|
||||
|
||||
from chatballs.ai.models import AIAgent, AIAgentStatus
|
||||
from chatballs.ai.provider.base import ProviderError
|
||||
from chatballs.channels.models import Channel
|
||||
from chatballs.conversations.models import (
|
||||
ConnectionIdentity,
|
||||
@@ -28,6 +27,7 @@ from chatballs.identity.models import (
|
||||
)
|
||||
from chatballs.integrations.models import Integration, IntegrationKind, IntegrationProvider
|
||||
from chatballs.notifications.models import Notification, NotificationAudience, NotificationType
|
||||
from chatballs.testing import ai_answer, ai_failure, run_pending_ai_turns
|
||||
from chatballs.testing import TenantAPIClient as APIClient
|
||||
|
||||
|
||||
@@ -65,13 +65,11 @@ class IngestProviderFailureTests(TestCase):
|
||||
from chatballs.conversations.ingest import ingest_inbound
|
||||
|
||||
with (
|
||||
mock.patch(
|
||||
"chatballs.conversations.ingest.run_channel_turn",
|
||||
side_effect=ProviderError("provider is down"),
|
||||
),
|
||||
mock.patch("chatballs.conversations.ingest.transports.send_reply", return_value=True) as send,
|
||||
ai_failure("provider is down"),
|
||||
mock.patch("chatballs.conversations.transports.send_reply", return_value=True) as send,
|
||||
):
|
||||
ingest_inbound(self.integration, self.inbound)
|
||||
run_pending_ai_turns()
|
||||
|
||||
conversation = self.channel.conversations.get()
|
||||
# Диалог передан оператору, ответчик — оператор.
|
||||
@@ -309,10 +307,11 @@ class ContactShareIngestTests(TestCase):
|
||||
external_id="ext-1", user_id="u1", chat_id="c1", text="Привет", display_name="Иван", username="ivan"
|
||||
)
|
||||
with (
|
||||
mock.patch("chatballs.conversations.ingest.run_channel_turn", return_value=mock.Mock(text="Здравствуйте!")),
|
||||
mock.patch("chatballs.conversations.ingest.transports.send_reply", return_value=True),
|
||||
ai_answer("Здравствуйте!"),
|
||||
mock.patch("chatballs.conversations.transports.send_reply", return_value=True),
|
||||
):
|
||||
ingest_inbound(self.integration, inbound)
|
||||
run_pending_ai_turns()
|
||||
|
||||
identity = ConnectionIdentity.objects.get(connection=self.integration, external_user_id="u1")
|
||||
self.assertEqual(identity.username, "ivan")
|
||||
@@ -323,13 +322,13 @@ class ContactShareIngestTests(TestCase):
|
||||
inbound = InboundMessage(
|
||||
external_id="ext-2", user_id="u1", chat_id="c1", text="", display_name="Иван", username="ivan", phone="+79991234567"
|
||||
)
|
||||
with (
|
||||
mock.patch("chatballs.conversations.ingest.run_channel_turn") as ai_turn,
|
||||
mock.patch("chatballs.conversations.ingest.transports.send_contact_ack", return_value=True) as ack,
|
||||
):
|
||||
with mock.patch(
|
||||
"chatballs.conversations.ingest.transports.send_contact_ack", return_value=True
|
||||
) as ack:
|
||||
ingest_inbound(self.integration, inbound)
|
||||
|
||||
ai_turn.assert_not_called()
|
||||
# Ход AI даже не заявлен: отвечать на присланный контакт нечего.
|
||||
self.assertEqual(run_pending_ai_turns(), 0)
|
||||
ack.assert_called_once()
|
||||
contact = ConnectionIdentity.objects.get(connection=self.integration, external_user_id="u1").contact
|
||||
self.assertEqual(contact.phone, "+79991234567")
|
||||
@@ -525,6 +524,7 @@ class WebchatContactTests(TestCase):
|
||||
)
|
||||
|
||||
response = self._post_message("Здравствуйте")
|
||||
run_pending_ai_turns()
|
||||
|
||||
self.assertEqual(response.status_code, 201)
|
||||
conversation = Conversation.objects.get(channel=self.channel)
|
||||
@@ -555,11 +555,9 @@ class WebchatContactTests(TestCase):
|
||||
)
|
||||
|
||||
def test_provider_error_hands_off_without_500(self) -> None:
|
||||
with mock.patch(
|
||||
"chatballs.conversations.ingest.run_channel_turn",
|
||||
side_effect=ProviderError("AI недоступен"),
|
||||
):
|
||||
with ai_failure("AI недоступен"):
|
||||
response = self._post_message("Здравствуйте")
|
||||
run_pending_ai_turns()
|
||||
|
||||
self.assertEqual(response.status_code, 201)
|
||||
conversation = Conversation.objects.get(channel=self.channel)
|
||||
|
||||
@@ -0,0 +1,81 @@
|
||||
"""Расшифровка голосовых сообщений через BYOK-провайдера организации.
|
||||
|
||||
Три шага (``prepare`` → ``run`` → ``store``) вместо одной функции: обращение к
|
||||
провайдеру ждёт ответа десятки секунд, и всё это время транзакция занимала бы
|
||||
соединение из пула (chatballs.tenancy.middleware). Кто может разнести шаги —
|
||||
разносит: ход AI (chatballs.conversations.ai_turn) и кнопка «расшифровать» в
|
||||
рабочем месте (chatballs.conversations.voice_views).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
|
||||
from chatballs.ai.provider.base import ProviderError
|
||||
from chatballs.conversations.models import Message, TranscriptStatus
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class TranscriptionJob:
|
||||
"""Всё, что нужно провайдеру, — уже прочитанное из базы и хранилища.
|
||||
|
||||
Разложено на три шага (``prepare`` → ``run`` → ``store``), чтобы вызывающий
|
||||
мог держать транзакцию только вокруг первого и третьего: обращение к
|
||||
провайдеру ждёт ответа десятки секунд, и всё это время транзакция занимала
|
||||
бы соединение из пула (chatballs.tenancy.middleware).
|
||||
"""
|
||||
|
||||
provider: object
|
||||
model: str
|
||||
audio: bytes
|
||||
filename: str
|
||||
content_type: str
|
||||
|
||||
|
||||
def prepare_transcription(channel, message: Message) -> TranscriptionJob | None:
|
||||
"""Шаг в транзакции: провайдер организации, модель и байты аудио."""
|
||||
from chatballs.ai.provider.factory import get_transcription_provider
|
||||
from chatballs.ai.provider.routing import (
|
||||
DEFAULT_TRANSCRIPTION_MODEL,
|
||||
resolve_transcription_model,
|
||||
)
|
||||
|
||||
if not message.audio:
|
||||
return None
|
||||
provider = get_transcription_provider(channel=channel)
|
||||
try:
|
||||
model = resolve_transcription_model(channel)
|
||||
except ProviderError:
|
||||
model = DEFAULT_TRANSCRIPTION_MODEL # тестовый провайдер без интеграции
|
||||
with message.audio.open("rb") as handle:
|
||||
audio = handle.read()
|
||||
return TranscriptionJob(
|
||||
provider=provider,
|
||||
model=model,
|
||||
audio=audio,
|
||||
filename=message.audio.name.rsplit("/", 1)[-1],
|
||||
content_type=message.audio_content_type or "audio/ogg",
|
||||
)
|
||||
|
||||
|
||||
def run_transcription(job: TranscriptionJob) -> str:
|
||||
"""Шаг без транзакции: обращение к провайдеру."""
|
||||
return job.provider.transcribe(
|
||||
audio=job.audio,
|
||||
filename=job.filename,
|
||||
content_type=job.content_type,
|
||||
model=job.model,
|
||||
).strip()
|
||||
|
||||
|
||||
def store_transcription(message: Message, transcript: str) -> None:
|
||||
"""Шаг в транзакции: сохранить стенограмму и статус."""
|
||||
message.transcript = transcript
|
||||
message.transcript_status = TranscriptStatus.READY if transcript else TranscriptStatus.FAILED
|
||||
message.save(update_fields=["transcript", "transcript_status"])
|
||||
|
||||
|
||||
def mark_transcription_failed(message: Message) -> None:
|
||||
"""Статус FAILED — оператор повторит кнопкой."""
|
||||
message.transcript_status = TranscriptStatus.FAILED
|
||||
message.save(update_fields=["transcript_status"])
|
||||
@@ -20,9 +20,12 @@ from dataclasses import dataclass
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Первая пауза — два цикла опроса, дальше удвоение до четверти часа.
|
||||
# Первая пауза — два цикла опроса, дальше удвоение до минуты. Потолок был
|
||||
# четвертью часа, пока опрос и ответы AI жили в одном процессе: длинная пауза
|
||||
# берегла общий цикл. Теперь опрос ничего не ждёт, а четверть часа тишины после
|
||||
# одного сетевого сбоя клиент видит как «бот молчит».
|
||||
FIRST_DELAY_SECONDS = 6.0
|
||||
MAX_DELAY_SECONDS = 900.0
|
||||
MAX_DELAY_SECONDS = 60.0
|
||||
|
||||
|
||||
@dataclass
|
||||
@@ -56,8 +59,8 @@ def record_failure(integration, error: object) -> None:
|
||||
)
|
||||
elif delay >= MAX_DELAY_SECONDS and (previous is None or previous.delay < MAX_DELAY_SECONDS):
|
||||
logger.warning(
|
||||
"%s poll keeps failing for integration %s: %s (retrying every %.0f min)",
|
||||
integration.provider, integration.id, error, MAX_DELAY_SECONDS / 60,
|
||||
"%s poll keeps failing for integration %s: %s (retrying every %.0fs)",
|
||||
integration.provider, integration.id, error, MAX_DELAY_SECONDS,
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -94,7 +94,10 @@ def _bot_started(update: dict) -> InboundMessage | None:
|
||||
|
||||
|
||||
def _normalize(update: dict) -> InboundMessage | None:
|
||||
logger.info("MAX raw update: %s", json.dumps(update, ensure_ascii=False))
|
||||
# Сырой апдейт нужен при разборе настройки, а не в каждой строке журнала
|
||||
# рабочего сервера: поля MAX документированы не полностью, и посмотреть их
|
||||
# глазами иногда надо — но по включённому DEBUG.
|
||||
logger.debug("MAX raw update: %s", json.dumps(update, ensure_ascii=False))
|
||||
update_type = update.get("update_type") or update.get("updateType")
|
||||
if update_type == "bot_started":
|
||||
return _bot_started(update)
|
||||
|
||||
@@ -75,7 +75,7 @@ class MessageTranscribeView(ConversationViewBase):
|
||||
tenant_manages_own_transaction = True
|
||||
|
||||
def post(self, request: Request, message_id: int) -> Response:
|
||||
from chatballs.conversations.ingest import (
|
||||
from chatballs.conversations.transcription import (
|
||||
mark_transcription_failed,
|
||||
prepare_transcription,
|
||||
run_transcription,
|
||||
|
||||
@@ -10,11 +10,28 @@ logger = logging.getLogger(__name__)
|
||||
|
||||
EventHandler = Callable[[dict, TenantContext | None], None]
|
||||
_REGISTRY: dict[str, EventHandler] = {}
|
||||
# Обработчики, которые открывают транзакции сами (см. `register`).
|
||||
_OWN_TRANSACTION: set[str] = set()
|
||||
|
||||
|
||||
def register(event_type: str) -> Callable[[EventHandler], EventHandler]:
|
||||
def register(
|
||||
event_type: str, *, manages_own_transaction: bool = False
|
||||
) -> Callable[[EventHandler], EventHandler]:
|
||||
"""Зарегистрировать обработчик события.
|
||||
|
||||
По умолчанию обработчик выполняется целиком в одной транзакции: так у него
|
||||
есть RLS-контекст и атомарность, и думать об этом не нужно. Обработчику,
|
||||
который ходит наружу — к модели, в мессенджер, — такая транзакция стоит
|
||||
соединения из пула на всё время ожидания. Он выставляет
|
||||
`manages_own_transaction` и открывает `tenant_atomic` сам, вокруг обращений
|
||||
к базе. Забытый блок не опасен: без транзакции RLS-настройка пуста и строки
|
||||
просто не видны — ошибка проявится сразу.
|
||||
"""
|
||||
|
||||
def decorator(handler: EventHandler) -> EventHandler:
|
||||
_REGISTRY[event_type] = handler
|
||||
if manages_own_transaction:
|
||||
_OWN_TRANSACTION.add(event_type)
|
||||
return handler
|
||||
|
||||
return decorator
|
||||
@@ -29,5 +46,8 @@ def dispatch(event: OutboxEvent) -> None:
|
||||
if context is None:
|
||||
handler(event.payload, None)
|
||||
return
|
||||
if event.event_type in _OWN_TRANSACTION:
|
||||
handler(event.payload, context)
|
||||
return
|
||||
with tenant_atomic(context):
|
||||
handler(event.payload, context)
|
||||
@@ -1,4 +1,5 @@
|
||||
import logging
|
||||
import os
|
||||
import time
|
||||
|
||||
from django.core.management import call_command
|
||||
@@ -11,7 +12,12 @@ from chatballs.conversations.maintenance import close_stale_conversations
|
||||
from chatballs.conversations.poller import poll_all_messengers
|
||||
from chatballs.events.handlers import dispatch
|
||||
from chatballs.events.models import OutboxStatus
|
||||
from chatballs.events.services import claim_next_outbox_event, mark_retry
|
||||
from chatballs.events.services import (
|
||||
OUTBOX_DB,
|
||||
claim_next_outbox_event,
|
||||
mark_retry,
|
||||
release_stale_processing,
|
||||
)
|
||||
from chatballs.notifications.binding import poll_notifier_bots
|
||||
from chatballs.tenancy.context import TenantActorKind, TenantContext
|
||||
from chatballs.tenancy.database import tenant_atomic
|
||||
@@ -20,16 +26,42 @@ from chatballs.updates.services import check_for_updates
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Роли воркера. Разделены потому, что работы у них разного веса: опрос
|
||||
# подключений — это короткие запросы и записи в базу, а обработка событий может
|
||||
# ждать модель десятки секунд. В одном процессе второе перекрывало первое, и
|
||||
# входящие переставали забираться на всё время ответа AI.
|
||||
ROLE_ALL = "all"
|
||||
ROLE_POLLER = "poller"
|
||||
ROLE_EVENTS = "events"
|
||||
|
||||
MESSENGER_POLL_INTERVAL = 3.0 # seconds between messenger long-poll cycles
|
||||
MAINTENANCE_INTERVAL = 3600.0 # seconds between maintenance cycles (auto-close stale dialogs)
|
||||
CALL_SWEEP_INTERVAL = 10.0 # seconds between call timeout sweeps (invite expiry, stuck connect)
|
||||
# Пороги очереди задаются в минутах, поэтому раз в полминуты — с запасом:
|
||||
# проверка дешёвая, а повтор гасится dedup-ключом уровня.
|
||||
QUEUE_SWEEP_INTERVAL = 30.0 # seconds between waiting-queue escalation sweeps
|
||||
# Возврат событий, взятых в работу упавшим процессом.
|
||||
STALE_SWEEP_INTERVAL = 60.0
|
||||
|
||||
|
||||
class Command(BaseCommand):
|
||||
help = "Runs the local domain event worker (outbox dispatch + messenger inbound polling)."
|
||||
help = (
|
||||
"Runs the local domain event worker. Roles: 'poller' polls messenger "
|
||||
"connections and runs sweeps, 'events' dispatches the outbox (AI turns), "
|
||||
"'all' does both in one process (development default)."
|
||||
)
|
||||
|
||||
def add_arguments(self, parser) -> None:
|
||||
parser.add_argument(
|
||||
"--role",
|
||||
choices=[ROLE_ALL, ROLE_POLLER, ROLE_EVENTS],
|
||||
default=os.environ.get("CHATBALLS_WORKER_ROLE", ROLE_ALL),
|
||||
help=(
|
||||
"Что делает этот процесс. Опрос держат в одном экземпляре "
|
||||
"(курсоры подключений и паузы после сбоя живут в его памяти), "
|
||||
"роль событий масштабируется репликами."
|
||||
),
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _tenant_contexts():
|
||||
@@ -38,35 +70,54 @@ class Command(BaseCommand):
|
||||
organization, actor_kind=TenantActorKind.SYSTEM
|
||||
)
|
||||
|
||||
def _for_each_tenant(self, operation, failure: str) -> None:
|
||||
try:
|
||||
for context in self._tenant_contexts():
|
||||
with tenant_atomic(context):
|
||||
operation(context)
|
||||
except Exception: # pragma: no cover
|
||||
logger.exception(failure)
|
||||
|
||||
def _dispatch_one(self) -> bool:
|
||||
"""Взять и обработать одно событие. False — событий нет."""
|
||||
|
||||
try:
|
||||
event = claim_next_outbox_event()
|
||||
except Exception: # pragma: no cover
|
||||
# Отравленное событие не должно ронять процесс: иначе воркер
|
||||
# уходит в краш-петлю и вместе с outbox встают поллинг
|
||||
# мессенджеров и таймауты звонков.
|
||||
logger.exception("Outbox claim cycle failed")
|
||||
time.sleep(1)
|
||||
return False
|
||||
if event is None:
|
||||
return False
|
||||
try:
|
||||
logger.info("Processing outbox event %s", event.id)
|
||||
dispatch(event)
|
||||
event.status = OutboxStatus.PROCESSED
|
||||
event.processed_at = timezone.now()
|
||||
event.save(
|
||||
using=OUTBOX_DB,
|
||||
update_fields=["status", "processed_at"],
|
||||
)
|
||||
except Exception as exc: # pragma: no cover
|
||||
logger.exception("Outbox event failed: %s", event.id)
|
||||
mark_retry(event, str(exc))
|
||||
return True
|
||||
|
||||
def handle(self, *args: object, **options: object) -> None:
|
||||
self.stdout.write("Hub worker started")
|
||||
role = str(options["role"])
|
||||
does_events = role in (ROLE_ALL, ROLE_EVENTS)
|
||||
does_polling = role in (ROLE_ALL, ROLE_POLLER)
|
||||
self.stdout.write(f"Hub worker started (role={role})")
|
||||
last_poll = 0.0
|
||||
last_maintenance = 0.0
|
||||
last_call_sweep = 0.0
|
||||
last_queue_sweep = 0.0
|
||||
last_stale_sweep = 0.0
|
||||
while True:
|
||||
try:
|
||||
event = claim_next_outbox_event()
|
||||
except Exception: # pragma: no cover
|
||||
# Отравленное событие не должно ронять процесс: иначе воркер
|
||||
# уходит в краш-петлю и вместе с outbox встают поллинг
|
||||
# мессенджеров и таймауты звонков.
|
||||
logger.exception("Outbox claim cycle failed")
|
||||
time.sleep(1)
|
||||
continue
|
||||
if event is not None:
|
||||
try:
|
||||
logger.info("Processing outbox event %s", event.id)
|
||||
dispatch(event)
|
||||
event.status = OutboxStatus.PROCESSED
|
||||
event.processed_at = timezone.now()
|
||||
event.save(
|
||||
using="platform",
|
||||
update_fields=["status", "processed_at"],
|
||||
)
|
||||
except Exception as exc: # pragma: no cover
|
||||
logger.exception("Outbox event failed: %s", event.id)
|
||||
mark_retry(event, str(exc))
|
||||
worked = self._dispatch_one() if does_events else False
|
||||
|
||||
# Дальше идут периодические работы. Раньше обработка события
|
||||
# обрывала цикл на `continue`, и при непрерывном потоке событий —
|
||||
@@ -74,37 +125,27 @@ class Command(BaseCommand):
|
||||
# переставали забираться входящие сообщения и истекать приглашения
|
||||
# на звонки. Проверки дешёвые: почти всегда это сравнение времени.
|
||||
now = time.monotonic()
|
||||
if now - last_poll >= MESSENGER_POLL_INTERVAL:
|
||||
if does_events and now - last_stale_sweep >= STALE_SWEEP_INTERVAL:
|
||||
last_stale_sweep = now
|
||||
try:
|
||||
released = release_stale_processing()
|
||||
if released:
|
||||
logger.warning("Released %s stale outbox event(s)", released)
|
||||
except Exception: # pragma: no cover
|
||||
logger.exception("Stale outbox sweep failed")
|
||||
if does_polling and now - last_poll >= MESSENGER_POLL_INTERVAL:
|
||||
last_poll = now
|
||||
try:
|
||||
for context in self._tenant_contexts():
|
||||
with tenant_atomic(context):
|
||||
poll_all_messengers(context)
|
||||
except Exception: # pragma: no cover
|
||||
logger.exception("Messenger polling cycle failed")
|
||||
try:
|
||||
for context in self._tenant_contexts():
|
||||
with tenant_atomic(context):
|
||||
poll_notifier_bots(context)
|
||||
except Exception: # pragma: no cover
|
||||
logger.exception("Notifier polling cycle failed")
|
||||
if now - last_call_sweep >= CALL_SWEEP_INTERVAL:
|
||||
self._for_each_tenant(poll_all_messengers, "Messenger polling cycle failed")
|
||||
self._for_each_tenant(poll_notifier_bots, "Notifier polling cycle failed")
|
||||
if does_polling and now - last_call_sweep >= CALL_SWEEP_INTERVAL:
|
||||
last_call_sweep = now
|
||||
try:
|
||||
for context in self._tenant_contexts():
|
||||
with tenant_atomic(context):
|
||||
expire_stale_calls(context)
|
||||
except Exception: # pragma: no cover
|
||||
logger.exception("Call sweep cycle failed")
|
||||
if now - last_queue_sweep >= QUEUE_SWEEP_INTERVAL:
|
||||
self._for_each_tenant(expire_stale_calls, "Call sweep cycle failed")
|
||||
if does_polling and now - last_queue_sweep >= QUEUE_SWEEP_INTERVAL:
|
||||
last_queue_sweep = now
|
||||
try:
|
||||
for context in self._tenant_contexts():
|
||||
with tenant_atomic(context):
|
||||
sweep_waiting_conversations(context)
|
||||
except Exception: # pragma: no cover
|
||||
logger.exception("Waiting queue sweep cycle failed")
|
||||
if now - last_maintenance >= MAINTENANCE_INTERVAL:
|
||||
self._for_each_tenant(
|
||||
sweep_waiting_conversations, "Waiting queue sweep cycle failed"
|
||||
)
|
||||
if does_polling and now - last_maintenance >= MAINTENANCE_INTERVAL:
|
||||
last_maintenance = now
|
||||
# Канал релизов спрашивается не чаще раза в несколько часов:
|
||||
# интервал держит сама проверка по времени последнего ответа.
|
||||
@@ -112,12 +153,7 @@ class Command(BaseCommand):
|
||||
check_for_updates()
|
||||
except Exception: # pragma: no cover
|
||||
logger.exception("Update check cycle failed")
|
||||
try:
|
||||
for context in self._tenant_contexts():
|
||||
with tenant_atomic(context):
|
||||
close_stale_conversations(context)
|
||||
except Exception: # pragma: no cover
|
||||
logger.exception("Maintenance cycle failed")
|
||||
self._for_each_tenant(close_stale_conversations, "Maintenance cycle failed")
|
||||
try:
|
||||
# Просроченные сессии Django сам не удаляет, а их накопление
|
||||
# утяжеляет карточку сотрудника: владельца сессии видно
|
||||
@@ -127,5 +163,5 @@ class Command(BaseCommand):
|
||||
logger.exception("Session cleanup failed")
|
||||
# Спим только когда работы нет: иначе очередь событий разбиралась бы
|
||||
# по одному событию в секунду.
|
||||
if event is None:
|
||||
if not worked:
|
||||
time.sleep(1)
|
||||
@@ -3,6 +3,7 @@ from datetime import timedelta
|
||||
from typing import Any
|
||||
|
||||
from django.db import transaction
|
||||
from django.db.models import Exists, OuterRef
|
||||
from django.utils import timezone
|
||||
|
||||
from chatballs.events.context import get_correlation_id
|
||||
@@ -12,6 +13,11 @@ from chatballs.tenancy.context import TenantActorKind, TenantContext
|
||||
from chatballs.tenancy.lookup import load_organization
|
||||
|
||||
|
||||
# Где живёт outbox. Захват идёт по всем организациям сразу, поэтому читает и
|
||||
# отмечает события роль platform, а не app (chatballs.tenancy.routing).
|
||||
OUTBOX_DB = "platform"
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class DomainEvent:
|
||||
aggregate_type: str
|
||||
@@ -92,24 +98,60 @@ def mark_retry(event: OutboxEvent, error: str, max_attempts: int = 5) -> None:
|
||||
event.status = OutboxStatus.DEAD_LETTER if event.attempts >= max_attempts else OutboxStatus.FAILED
|
||||
event.next_attempt_at = timezone.now() + timedelta(seconds=min(300, 2**event.attempts))
|
||||
event.save(
|
||||
using="platform",
|
||||
using=OUTBOX_DB,
|
||||
update_fields=["attempts", "last_error", "status", "next_attempt_at"],
|
||||
)
|
||||
|
||||
|
||||
def claim_next_outbox_event() -> OutboxEvent | None:
|
||||
with transaction.atomic(using="platform"):
|
||||
# Сколько событию отведено на обработку. Роль событий работает в нескольких
|
||||
# процессах, и взятое в работу событие не должно достаться второму; но и
|
||||
# пропасть навсегда, если процесс упал посреди обработки, оно тоже не должно.
|
||||
# Срок хранится в `next_attempt_at`: у поля ровно этот смысл — «не раньше».
|
||||
PROCESSING_LEASE_SECONDS = 300
|
||||
|
||||
|
||||
def claim_next_outbox_event(*, lease_seconds: int = PROCESSING_LEASE_SECONDS) -> OutboxEvent | None:
|
||||
"""Взять следующее событие в работу.
|
||||
|
||||
События одного объекта идут строго по очереди: пока по агрегату есть
|
||||
событие в работе, следующее не выдаётся. Иначе два ответа AI одному
|
||||
диалогу считались бы параллельно и приезжали клиенту вперемешку.
|
||||
"""
|
||||
now = timezone.now()
|
||||
busy = OutboxEvent.objects.using(OUTBOX_DB).filter(
|
||||
status=OutboxStatus.PROCESSING,
|
||||
aggregate_type=OuterRef("aggregate_type"),
|
||||
aggregate_id=OuterRef("aggregate_id"),
|
||||
)
|
||||
with transaction.atomic(using=OUTBOX_DB):
|
||||
event = (
|
||||
OutboxEvent.objects.using("platform").select_for_update(skip_locked=True)
|
||||
OutboxEvent.objects.using(OUTBOX_DB).select_for_update(skip_locked=True)
|
||||
.filter(
|
||||
status__in=[OutboxStatus.PENDING, OutboxStatus.FAILED],
|
||||
next_attempt_at__lte=timezone.now(),
|
||||
next_attempt_at__lte=now,
|
||||
)
|
||||
.filter(~Exists(busy))
|
||||
.order_by("next_attempt_at", "created_at")
|
||||
.first()
|
||||
)
|
||||
if event is None:
|
||||
return None
|
||||
event.status = OutboxStatus.PROCESSING
|
||||
event.save(using="platform", update_fields=["status"])
|
||||
event.next_attempt_at = now + timedelta(seconds=lease_seconds)
|
||||
event.save(using=OUTBOX_DB, update_fields=["status", "next_attempt_at"])
|
||||
return event
|
||||
|
||||
|
||||
def release_stale_processing() -> int:
|
||||
"""Вернуть в очередь события, взятые в работу и не доведённые до конца.
|
||||
|
||||
Процесс мог упасть или его перезапустили между `claim` и записью
|
||||
результата. Без возврата такое событие остаётся `PROCESSING` навсегда — а
|
||||
вместе с ним встаёт и весь агрегат, потому что следующие события того же
|
||||
объекта ждут его.
|
||||
"""
|
||||
return (
|
||||
OutboxEvent.objects.using(OUTBOX_DB)
|
||||
.filter(status=OutboxStatus.PROCESSING, next_attempt_at__lte=timezone.now())
|
||||
.update(status=OutboxStatus.PENDING)
|
||||
)
|
||||
@@ -0,0 +1,75 @@
|
||||
"""Очередь событий: порядок внутри агрегата и возврат зависших."""
|
||||
|
||||
from datetime import timedelta
|
||||
from unittest import mock
|
||||
|
||||
from django.test import TestCase
|
||||
from django.utils import timezone
|
||||
|
||||
from chatballs.events.models import EventOwnership, OutboxEvent, OutboxStatus
|
||||
from chatballs.events.services import claim_next_outbox_event, release_stale_processing
|
||||
from chatballs.identity.bootstrap import bootstrap_owner
|
||||
from chatballs.identity.models import Organization
|
||||
from chatballs.tenancy.context import TenantActorKind
|
||||
|
||||
|
||||
class OutboxClaimTests(TestCase):
|
||||
"""Захват событий.
|
||||
|
||||
В проде outbox читает роль platform, в тестах отдельного соединения под неё
|
||||
нет — все алиасы смотрят в одну базу. Подменяется только алиас: сами
|
||||
запросы те же, что в проде.
|
||||
"""
|
||||
|
||||
def setUp(self) -> None:
|
||||
bootstrap_owner(email="owner@example.com", password="temporary-password")
|
||||
self.organization = Organization.objects.get(slug="demo")
|
||||
patch = mock.patch("chatballs.events.services.OUTBOX_DB", "default")
|
||||
patch.start()
|
||||
self.addCleanup(patch.stop)
|
||||
|
||||
def _event(self, aggregate_id: str, event_type: str = "conversation.ai_turn_requested"):
|
||||
return OutboxEvent.objects.create(
|
||||
aggregate_type="Conversation",
|
||||
aggregate_id=aggregate_id,
|
||||
event_type=event_type,
|
||||
payload={},
|
||||
ownership=EventOwnership.TENANT,
|
||||
organization=self.organization,
|
||||
actor_kind=TenantActorKind.MACHINE,
|
||||
)
|
||||
|
||||
def test_second_event_of_the_same_aggregate_waits(self) -> None:
|
||||
# Два ответа одному диалогу не считаются параллельно: иначе они
|
||||
# приезжают клиенту вперемешку.
|
||||
self._event("7")
|
||||
self._event("7")
|
||||
|
||||
first = claim_next_outbox_event()
|
||||
self.assertIsNotNone(first)
|
||||
self.assertIsNone(claim_next_outbox_event())
|
||||
|
||||
def test_other_aggregates_are_not_blocked(self) -> None:
|
||||
self._event("7")
|
||||
self._event("8")
|
||||
|
||||
self.assertIsNotNone(claim_next_outbox_event())
|
||||
self.assertIsNotNone(claim_next_outbox_event())
|
||||
|
||||
def test_stale_processing_returns_to_the_queue(self) -> None:
|
||||
# Процесс упал между взятием события и записью результата: без
|
||||
# возврата оно осталось бы в работе навсегда, а вместе с ним встал бы
|
||||
# весь диалог.
|
||||
event = self._event("7")
|
||||
claimed = claim_next_outbox_event()
|
||||
self.assertEqual(claimed.id, event.id)
|
||||
|
||||
self.assertEqual(release_stale_processing(), 0)
|
||||
|
||||
OutboxEvent.objects.filter(id=event.id).update(
|
||||
next_attempt_at=timezone.now() - timedelta(seconds=1)
|
||||
)
|
||||
self.assertEqual(release_stale_processing(), 1)
|
||||
event.refresh_from_db()
|
||||
self.assertEqual(event.status, OutboxStatus.PENDING)
|
||||
self.assertIsNotNone(claim_next_outbox_event())
|
||||
@@ -58,6 +58,7 @@ MESSAGES: dict[str, object] = {
|
||||
"conversations.system.assigned_to": "Conversation assigned to {operator}",
|
||||
"conversations.system.assignment_expired": "{operator} did not pick the conversation up — it is back in the queue",
|
||||
"conversations.system.ai_handed_over": "AI handed the conversation to an operator",
|
||||
"conversations.ai_unavailable_reply": "Sorry, I cannot answer right now. I have passed your question to a specialist — they will join shortly.",
|
||||
"conversations.system.ai_unavailable": "AI is unavailable — the conversation was handed to an operator",
|
||||
"conversations.system.call_accepted": "The customer accepted the call invitation",
|
||||
"conversations.system.call_cancelled": "The operator cancelled the call invitation",
|
||||
|
||||
@@ -62,6 +62,7 @@ MESSAGES: dict[str, object] = {
|
||||
"conversations.system.ai_handed_over": "AI передал диалог оператору",
|
||||
"conversations.system.assigned_to": "Диалог назначен на {operator}",
|
||||
"conversations.system.assignment_expired": "{operator} не взял диалог — он вернулся в очередь",
|
||||
"conversations.ai_unavailable_reply": "Извините, прямо сейчас не получается ответить. Я передал ваш вопрос специалисту — он скоро подключится.",
|
||||
"conversations.system.ai_unavailable": "AI недоступен — диалог передан оператору",
|
||||
"conversations.system.call_accepted": "Клиент принял приглашение на звонок",
|
||||
"conversations.system.call_cancelled": "Сотрудник отменил приглашение на звонок",
|
||||
|
||||
@@ -23,6 +23,57 @@ def system_tenant_context(organization: Organization) -> TenantContext:
|
||||
return TenantContext.for_resource(organization, actor_kind=TenantActorKind.SYSTEM)
|
||||
|
||||
|
||||
def run_pending_ai_turns() -> int:
|
||||
"""Прогнать поставленные ходы AI и вернуть их число.
|
||||
|
||||
Ход считается ролью событий воркера, которой в тестах нет: приём ставит
|
||||
заявку, а вызвать её должен сам тест — так же, как он это делает с
|
||||
приглашениями на звонок и письмами.
|
||||
"""
|
||||
from chatballs.conversations.ai_turn import AI_TURN_REQUESTED
|
||||
from chatballs.events.handlers import dispatch
|
||||
from chatballs.events.models import OutboxEvent, OutboxStatus
|
||||
|
||||
events = list(
|
||||
OutboxEvent.objects.filter(
|
||||
event_type=AI_TURN_REQUESTED, status=OutboxStatus.PENDING
|
||||
).order_by("created_at")
|
||||
)
|
||||
for event in events:
|
||||
dispatch(event)
|
||||
event.status = OutboxStatus.PROCESSED
|
||||
event.save(update_fields=["status"])
|
||||
return len(events)
|
||||
|
||||
|
||||
def ai_answer(text: str):
|
||||
"""Подменить ответ модели в ходе AI (контекст-менеджер)."""
|
||||
from unittest import mock
|
||||
|
||||
from chatballs.ai.provider.base import ChatResult
|
||||
from chatballs.ai.turn import TurnAnswer
|
||||
|
||||
return mock.patch(
|
||||
"chatballs.conversations.ai_turn.run_turn_chat",
|
||||
return_value=TurnAnswer(
|
||||
result=ChatResult(text=text, model="test", prompt_tokens=1, completion_tokens=1)
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def ai_failure(message: str = "provider is down"):
|
||||
"""Подменить ход AI отказом провайдера (контекст-менеджер)."""
|
||||
from unittest import mock
|
||||
|
||||
from chatballs.ai.provider.base import ProviderError
|
||||
from chatballs.ai.turn import TurnAnswer
|
||||
|
||||
return mock.patch(
|
||||
"chatballs.conversations.ai_turn.run_turn_chat",
|
||||
return_value=TurnAnswer(error=ProviderError(message)),
|
||||
)
|
||||
|
||||
|
||||
class TenantAPIClient(APIClient):
|
||||
"""Test client that turns legacy test literals into the C03 tenant route.
|
||||
|
||||
|
||||
@@ -8,6 +8,7 @@ from django.conf import settings
|
||||
from django.db import models, transaction
|
||||
from django.utils import timezone
|
||||
|
||||
from chatballs.conversations.ai_turn import conversation_is_thinking
|
||||
from chatballs.conversations.ingest import ingest_inbound
|
||||
from chatballs.conversations.models import (
|
||||
ConnectionIdentity,
|
||||
@@ -346,12 +347,15 @@ def messages_payload(session: WebSession, since: int) -> dict:
|
||||
"messages": [],
|
||||
"call": None,
|
||||
"reset": reset,
|
||||
"thinking": False,
|
||||
}
|
||||
items = conversation.messages.filter(id__gt=since).order_by("created_at")
|
||||
return {
|
||||
"reset": reset,
|
||||
"state": _STATE.get(conversation.control_mode, "ai"),
|
||||
"lifecycle": conversation.lifecycle,
|
||||
# Ответ уже считается: виджет показывает клиенту, что агент печатает.
|
||||
"thinking": conversation_is_thinking(conversation.id),
|
||||
"messages": [
|
||||
{
|
||||
"id": m.id,
|
||||
|
||||
@@ -197,15 +197,24 @@ CHATBALLS_OPENROUTER_BASE_URL = os.environ.get("CHATBALLS_OPENROUTER_BASE_URL",
|
||||
CHATBALLS_AI_REQUEST_TIMEOUT = float(os.environ.get("CHATBALLS_AI_REQUEST_TIMEOUT", "30"))
|
||||
CHATBALLS_AI_MAX_RETRIES = int(os.environ.get("CHATBALLS_AI_MAX_RETRIES", "2"))
|
||||
CHATBALLS_AI_EMBEDDING_MODEL = os.environ.get("CHATBALLS_AI_EMBEDDING_MODEL", "openai/text-embedding-3-small")
|
||||
# Ответ живому человеку ждут иначе, чем индексацию знаний: клиент в чате не
|
||||
# станет ждать полторы минуты, пока провайдер доберёт три попытки по тридцать
|
||||
# секунд. Ход диалога ограничен своим сроком (chatballs.ai.turn).
|
||||
CHATBALLS_AI_TURN_TIMEOUT = float(os.environ.get("CHATBALLS_AI_TURN_TIMEOUT", "20"))
|
||||
# Сколько ход имеет смысл: сообщение, пролежавшее в очереди дольше, отвечать
|
||||
# уже поздно — клиент ушёл или его взял оператор, и диалог передаётся человеку.
|
||||
CHATBALLS_AI_TURN_DEADLINE_SECONDS = int(
|
||||
os.environ.get("CHATBALLS_AI_TURN_DEADLINE_SECONDS", "120")
|
||||
)
|
||||
# Модель расшифровки голосовых (OpenAI-совместимый /audio/transcriptions).
|
||||
|
||||
# Managed-провайдер CustoAI удалён (ADR-CHATBALLS-0042 §3): AI — только через
|
||||
# интеграцию организации (BYOK).
|
||||
|
||||
# Long-poll hold-time мессенджеров (сек). Держим малым: единый воркер выполняет
|
||||
# и inbound-поллинг, и outbox-диспатч в одном потоке — при большом hold-time
|
||||
# getUpdates/updates блокирует цикл и outbox (приглашения звонков, уведомления,
|
||||
# ответы AI) уходит с задержкой в размер long-poll на каждое подключение.
|
||||
# Long-poll hold-time мессенджеров (сек). Держим малым: опрос идёт по
|
||||
# подключениям последовательно в одном процессе, и hold-time каждого из них
|
||||
# складывается в задержку приёма у остальных. Ответы AI от этого больше не
|
||||
# зависят — их считает отдельная роль воркера (run_worker --role).
|
||||
CHATBALLS_MESSENGER_POLL_TIMEOUT_SECONDS = int(os.environ.get("CHATBALLS_MESSENGER_POLL_TIMEOUT_SECONDS", "2"))
|
||||
|
||||
# Срок жизни анонимной сессии виджета: отсчёт от последней активности, а не от
|
||||
|
||||
@@ -116,6 +116,17 @@ services:
|
||||
- ./apps/backend:/app/apps/backend
|
||||
- ./data/media:/app/apps/backend/media
|
||||
|
||||
worker-events:
|
||||
build:
|
||||
context: .
|
||||
dockerfile: apps/backend/Dockerfile
|
||||
environment:
|
||||
CHATBALLS_DEBUG: "true"
|
||||
CHATBALLS_DELIVERY_MODE: ${CHATBALLS_DELIVERY_MODE:-CLOUD}
|
||||
volumes:
|
||||
- ./apps/backend:/app/apps/backend
|
||||
- ./data/media:/app/apps/backend/media
|
||||
|
||||
# Dev: образ сервиса обновления собирается локально; сам сервис в dev
|
||||
# бесполезен (стек поднят из исходников), но должен собираться и стартовать.
|
||||
updater:
|
||||
|
||||
+27
-1
@@ -200,10 +200,13 @@ services:
|
||||
redis:
|
||||
condition: service_healthy
|
||||
|
||||
# Опрос подключений и периодические работы. Строго один экземпляр: курсоры
|
||||
# мессенджеров и паузы после сбоя живут в памяти процесса, и второй опросчик
|
||||
# забирал бы те же обновления второй раз.
|
||||
worker:
|
||||
image: ${CHATBALLS_BACKEND_IMAGE:-chatballs-backend:dev}
|
||||
restart: unless-stopped
|
||||
command: ["python", "manage.py", "run_worker"]
|
||||
command: ["python", "manage.py", "run_worker", "--role=poller"]
|
||||
environment:
|
||||
CHATBALLS_DB_ROLE: app
|
||||
# Захват outbox идёт по всем организациям сразу — воркеру нужен алиас
|
||||
@@ -217,6 +220,29 @@ services:
|
||||
backend-app:
|
||||
condition: service_healthy
|
||||
|
||||
# Обработка событий: ответы AI, доставка уведомлений, приглашения на звонки.
|
||||
# Ход AI ждёт модель секунды и десятки секунд, поэтому эта работа вынесена из
|
||||
# опроса и масштабируется репликами — события разбираются с `skip_locked`, а
|
||||
# ходы одного диалога всё равно идут по очереди
|
||||
# (chatballs.events.services.claim_next_outbox_event). Реплик стоит держать
|
||||
# немного: они ходят к провайдеру одним ключом организации.
|
||||
worker-events:
|
||||
image: ${CHATBALLS_BACKEND_IMAGE:-chatballs-backend:dev}
|
||||
restart: unless-stopped
|
||||
command: ["python", "manage.py", "run_worker", "--role=events"]
|
||||
deploy:
|
||||
replicas: ${CHATBALLS_EVENT_WORKERS:-2}
|
||||
environment:
|
||||
CHATBALLS_DB_ROLE: app
|
||||
CHATBALLS_DB_PLATFORM_ALIAS: "1"
|
||||
volumes:
|
||||
- chatballs-secrets:/run/chatballs/secrets:ro
|
||||
- chatballs-secrets-platform:/run/chatballs/secrets/platform:ro
|
||||
- chatballs-media:/app/apps/backend/media
|
||||
depends_on:
|
||||
backend-app:
|
||||
condition: service_healthy
|
||||
|
||||
# Обновление по кнопке из интерфейса (ADR-CHATBALLS-0049). Единственный
|
||||
# сервис с доступом к docker.sock: забирает запрос из тома chatballs-updates,
|
||||
# проверяет, что это официальный релиз (адрес со страницы релизов, образы по
|
||||
|
||||
Reference in new issue
Block a user