diff --git a/apps/backend/chatballs/ai/invocation.py b/apps/backend/chatballs/ai/invocation.py index 5bdd0fb..764f984 100644 --- a/apps/backend/chatballs/ai/invocation.py +++ b/apps/backend/chatballs/ai/invocation.py @@ -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 diff --git a/apps/backend/chatballs/ai/provider/base.py b/apps/backend/chatballs/ai/provider/base.py index 153b0b0..3dfcdca 100644 --- a/apps/backend/chatballs/ai/provider/base.py +++ b/apps/backend/chatballs/ai/provider/base.py @@ -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" diff --git a/apps/backend/chatballs/ai/provider/factory.py b/apps/backend/chatballs/ai/provider/factory.py index 8087213..11a8e0c 100644 --- a/apps/backend/chatballs/ai/provider/factory.py +++ b/apps/backend/chatballs/ai/provider/factory.py @@ -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) diff --git a/apps/backend/chatballs/ai/provider/openai_http.py b/apps/backend/chatballs/ai/provider/openai_http.py index 40fef0f..77d4a4a 100644 --- a/apps/backend/chatballs/ai/provider/openai_http.py +++ b/apps/backend/chatballs/ai/provider/openai_http.py @@ -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. diff --git a/apps/backend/chatballs/ai/provider/resilience.py b/apps/backend/chatballs/ai/provider/resilience.py index 704da81..ce92918 100644 --- a/apps/backend/chatballs/ai/provider/resilience.py +++ b/apps/backend/chatballs/ai/provider/resilience.py @@ -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() diff --git a/apps/backend/chatballs/ai/provider/routing.py b/apps/backend/chatballs/ai/provider/routing.py index 94a49dc..7a5a150 100644 --- a/apps/backend/chatballs/ai/provider/routing.py +++ b/apps/backend/chatballs/ai/provider/routing.py @@ -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( diff --git a/apps/backend/chatballs/ai/retrieval.py b/apps/backend/chatballs/ai/retrieval.py index 76b86ed..f464876 100644 --- a/apps/backend/chatballs/ai/retrieval.py +++ b/apps/backend/chatballs/ai/retrieval.py @@ -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) diff --git a/apps/backend/chatballs/ai/runtime.py b/apps/backend/chatballs/ai/runtime.py index cd13fce..7a13bd3 100644 --- a/apps/backend/chatballs/ai/runtime.py +++ b/apps/backend/chatballs/ai/runtime.py @@ -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, diff --git a/apps/backend/chatballs/ai/turn.py b/apps/backend/chatballs/ai/turn.py new file mode 100644 index 0000000..f161f22 --- /dev/null +++ b/apps/backend/chatballs/ai/turn.py @@ -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, + ) diff --git a/apps/backend/chatballs/conversations/ai_turn.py b/apps/backend/chatballs/conversations/ai_turn.py new file mode 100644 index 0000000..56223e8 --- /dev/null +++ b/apps/backend/chatballs/conversations/ai_turn.py @@ -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) diff --git a/apps/backend/chatballs/conversations/ai_turn_result.py b/apps/backend/chatballs/conversations/ai_turn_result.py new file mode 100644 index 0000000..d2ab19d --- /dev/null +++ b/apps/backend/chatballs/conversations/ai_turn_result.py @@ -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 diff --git a/apps/backend/chatballs/conversations/apps.py b/apps/backend/chatballs/conversations/apps.py index 5de9d0a..a6d0b6f 100644 --- a/apps/backend/chatballs/conversations/apps.py +++ b/apps/backend/chatballs/conversations/apps.py @@ -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) diff --git a/apps/backend/chatballs/conversations/event_handlers.py b/apps/backend/chatballs/conversations/event_handlers.py new file mode 100644 index 0000000..9254442 --- /dev/null +++ b/apps/backend/chatballs/conversations/event_handlers.py @@ -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) diff --git a/apps/backend/chatballs/conversations/ingest.py b/apps/backend/chatballs/conversations/ingest.py index 1c3e156..1f97104 100644 --- a/apps/backend/chatballs/conversations/ingest.py +++ b/apps/backend/chatballs/conversations/ingest.py @@ -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: diff --git a/apps/backend/chatballs/conversations/migrations/0026_message_ai_turn_state.py b/apps/backend/chatballs/conversations/migrations/0026_message_ai_turn_state.py new file mode 100644 index 0000000..e772fb6 --- /dev/null +++ b/apps/backend/chatballs/conversations/migrations/0026_message_ai_turn_state.py @@ -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), + ), + ] diff --git a/apps/backend/chatballs/conversations/models.py b/apps/backend/chatballs/conversations/models.py index 6c6ceda..28404fc 100644 --- a/apps/backend/chatballs/conversations/models.py +++ b/apps/backend/chatballs/conversations/models.py @@ -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) diff --git a/apps/backend/chatballs/conversations/test_ai_turn.py b/apps/backend/chatballs/conversations/test_ai_turn.py new file mode 100644 index 0000000..c630d9e --- /dev/null +++ b/apps/backend/chatballs/conversations/test_ai_turn.py @@ -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() diff --git a/apps/backend/chatballs/conversations/test_email_transport.py b/apps/backend/chatballs/conversations/test_email_transport.py index 50b78df..f7368ee 100644 --- a/apps/backend/chatballs/conversations/test_email_transport.py +++ b/apps/backend/chatballs/conversations/test_email_transport.py @@ -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: diff --git a/apps/backend/chatballs/conversations/test_lifecycle.py b/apps/backend/chatballs/conversations/test_lifecycle.py index c61121d..d6db42f 100644 --- a/apps/backend/chatballs/conversations/test_lifecycle.py +++ b/apps/backend/chatballs/conversations/test_lifecycle.py @@ -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() diff --git a/apps/backend/chatballs/conversations/test_queue_order.py b/apps/backend/chatballs/conversations/test_queue_order.py index 97e57d3..0449d0f 100644 --- a/apps/backend/chatballs/conversations/test_queue_order.py +++ b/apps/backend/chatballs/conversations/test_queue_order.py @@ -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( diff --git a/apps/backend/chatballs/conversations/test_voice_ai_and_features.py b/apps/backend/chatballs/conversations/test_voice_ai_and_features.py index 98a5ed5..8fd3119 100644 --- a/apps/backend/chatballs/conversations/test_voice_ai_and_features.py +++ b/apps/backend/chatballs/conversations/test_voice_ai_and_features.py @@ -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): diff --git a/apps/backend/chatballs/conversations/tests.py b/apps/backend/chatballs/conversations/tests.py index f5e0d56..d888f97 100644 --- a/apps/backend/chatballs/conversations/tests.py +++ b/apps/backend/chatballs/conversations/tests.py @@ -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) diff --git a/apps/backend/chatballs/conversations/transcription.py b/apps/backend/chatballs/conversations/transcription.py new file mode 100644 index 0000000..892ff79 --- /dev/null +++ b/apps/backend/chatballs/conversations/transcription.py @@ -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"]) diff --git a/apps/backend/chatballs/conversations/transports/backoff.py b/apps/backend/chatballs/conversations/transports/backoff.py index fa26d96..ac96946 100644 --- a/apps/backend/chatballs/conversations/transports/backoff.py +++ b/apps/backend/chatballs/conversations/transports/backoff.py @@ -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, ) diff --git a/apps/backend/chatballs/conversations/transports/max.py b/apps/backend/chatballs/conversations/transports/max.py index dfabc65..8de4d59 100644 --- a/apps/backend/chatballs/conversations/transports/max.py +++ b/apps/backend/chatballs/conversations/transports/max.py @@ -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) diff --git a/apps/backend/chatballs/conversations/voice_views.py b/apps/backend/chatballs/conversations/voice_views.py index a4c6d55..1213eb1 100644 --- a/apps/backend/chatballs/conversations/voice_views.py +++ b/apps/backend/chatballs/conversations/voice_views.py @@ -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, diff --git a/apps/backend/chatballs/events/handlers.py b/apps/backend/chatballs/events/handlers.py index 0ff8c0b..6c7952d 100644 --- a/apps/backend/chatballs/events/handlers.py +++ b/apps/backend/chatballs/events/handlers.py @@ -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) diff --git a/apps/backend/chatballs/events/management/commands/run_worker.py b/apps/backend/chatballs/events/management/commands/run_worker.py index 07844e6..cd991e5 100644 --- a/apps/backend/chatballs/events/management/commands/run_worker.py +++ b/apps/backend/chatballs/events/management/commands/run_worker.py @@ -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) diff --git a/apps/backend/chatballs/events/services.py b/apps/backend/chatballs/events/services.py index 80ec17b..c643801 100644 --- a/apps/backend/chatballs/events/services.py +++ b/apps/backend/chatballs/events/services.py @@ -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) + ) diff --git a/apps/backend/chatballs/events/test_outbox_claim.py b/apps/backend/chatballs/events/test_outbox_claim.py new file mode 100644 index 0000000..42e9f5c --- /dev/null +++ b/apps/backend/chatballs/events/test_outbox_claim.py @@ -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()) diff --git a/apps/backend/chatballs/i18n/messages/en.py b/apps/backend/chatballs/i18n/messages/en.py index 6d44efe..14d3b33 100644 --- a/apps/backend/chatballs/i18n/messages/en.py +++ b/apps/backend/chatballs/i18n/messages/en.py @@ -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", diff --git a/apps/backend/chatballs/i18n/messages/ru.py b/apps/backend/chatballs/i18n/messages/ru.py index 1852f9f..ddd74d3 100644 --- a/apps/backend/chatballs/i18n/messages/ru.py +++ b/apps/backend/chatballs/i18n/messages/ru.py @@ -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": "Сотрудник отменил приглашение на звонок", diff --git a/apps/backend/chatballs/testing.py b/apps/backend/chatballs/testing.py index 9bc6e77..492bf05 100644 --- a/apps/backend/chatballs/testing.py +++ b/apps/backend/chatballs/testing.py @@ -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. diff --git a/apps/backend/chatballs/webchat/services.py b/apps/backend/chatballs/webchat/services.py index ceed676..4e36b3d 100644 --- a/apps/backend/chatballs/webchat/services.py +++ b/apps/backend/chatballs/webchat/services.py @@ -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, diff --git a/apps/backend/chatballs_backend/settings_base.py b/apps/backend/chatballs_backend/settings_base.py index 8b33044..a2e51eb 100644 --- a/apps/backend/chatballs_backend/settings_base.py +++ b/apps/backend/chatballs_backend/settings_base.py @@ -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")) # Срок жизни анонимной сессии виджета: отсчёт от последней активности, а не от diff --git a/compose.dev.yaml b/compose.dev.yaml index e55414b..7ec2320 100644 --- a/compose.dev.yaml +++ b/compose.dev.yaml @@ -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: diff --git a/compose.yaml b/compose.yaml index df4dc86..fbd14e2 100644 --- a/compose.yaml +++ b/compose.yaml @@ -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, # проверяет, что это официальный релиз (адрес со страницы релизов, образы по