mirror of
https://github.com/dartdavros/chatballs.git
synced 2026-10-05 09:14:58 +03:00
Кнопка «Завершить» у оператора обязана заканчивать звонок из любой живой фазы. Отмена умела только REQUESTED/RINGING, и после того как клиент принял вызов, а соединение не установилось (обычное дело за NAT), оператор получал 409 и звонок висел. Фазу выбирает finish_call — тот же код, что и на стороне клиента; в интерфейсе кнопка больше не молчит при ошибке, а заканчивает звонок вторым путём. Фото контакта из Telegram и MAX скачивается и хранится у нас, а отдаётся своим адресом: страница рабочего места живёт под CSP «img-src 'self'», и ссылка на CDN мессенджера до экрана не доезжала — оператор видел инициалы. Источник запоминается, поэтому фото качается один раз; у Telegram оно спрашивается отдельным запросом, которого в апдейте нет. Голосовое из MAX с незнакомой формой вложения больше не пропадает: раньше такое сообщение уходило в никуда, теперь оператор видит его заглушкой, а в журнал попадает сам payload. В журнал же пишется причина, по которой диалог сразу уходит в очередь: у канала нет активного AI-агента. Проверено: 83 теста звонков, 6 новых тестов фото контакта, тесты ingest, вложений и голосовых, typecheck рабочего места. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
314 lines
12 KiB
Python
314 lines
12 KiB
Python
from __future__ import annotations
|
||
|
||
from dataclasses import dataclass
|
||
from datetime import timedelta
|
||
|
||
from django.conf import settings
|
||
from django.db import IntegrityError, transaction
|
||
from django.utils import timezone
|
||
|
||
from chatballs.calls.errors import (
|
||
CallAccessDenied,
|
||
CallConflict,
|
||
CallInvalidTransition,
|
||
CallTokenError,
|
||
)
|
||
from chatballs.calls.lifecycle import finish_call, transition_call
|
||
from chatballs.calls.metrics import record_call_metric
|
||
from chatballs.calls.models import (
|
||
TERMINAL_CALL_STATUSES,
|
||
UNFINISHED_CALL_STATUSES,
|
||
CallEndedBy,
|
||
CallInvite,
|
||
CallKind,
|
||
CallParticipant,
|
||
CallSession,
|
||
CallStatus,
|
||
InviteDeliveryStatus,
|
||
ParticipantSide,
|
||
)
|
||
from chatballs.calls.permissions import ensure_call_access, ensure_conversation_call_access
|
||
from chatballs.calls.public_access import (
|
||
ResolvedInvite,
|
||
accept_call_by_access_token,
|
||
authorize_call_access_context,
|
||
authorize_call_access_token,
|
||
call_state_by_access_token,
|
||
decline_call_by_access_token,
|
||
end_call_by_access_token,
|
||
resolve_invite,
|
||
)
|
||
from chatballs.calls.tokens import (
|
||
issue_call_access_token,
|
||
issue_invite_token,
|
||
)
|
||
from chatballs.conversations.models import (
|
||
ConnectionIdentity,
|
||
ControlMode,
|
||
Conversation,
|
||
LifecycleState,
|
||
Message,
|
||
MessageAuthor,
|
||
SystemEvent,
|
||
)
|
||
from chatballs.conversations.services import ClaimError, claim_locked_conversation
|
||
from chatballs.events.services import DomainEvent, enqueue_event
|
||
from chatballs.i18n import t
|
||
from chatballs.integrations.features import call_allowed
|
||
from chatballs.integrations.models import IntegrationProvider
|
||
from chatballs.tenancy.context import TenantContext
|
||
|
||
__all__ = (
|
||
"accept_call_by_access_token",
|
||
"authorize_call_access_context",
|
||
"authorize_call_access_token",
|
||
"call_state_by_access_token",
|
||
"decline_call_by_access_token",
|
||
"end_call_by_access_token",
|
||
"resolve_invite",
|
||
"record_call_metric",
|
||
)
|
||
|
||
# Outbox-событие доставки приглашения в TG/MAX (обработчик — calls.event_handlers).
|
||
CALL_INVITE_SEND = "calls.invite_send"
|
||
|
||
|
||
@dataclass(frozen=True)
|
||
class CreatedCall:
|
||
call_session: CallSession
|
||
invite_token: str
|
||
staff_access_token: str
|
||
|
||
|
||
def _conversation_identity(conversation: Conversation) -> ConnectionIdentity:
|
||
if conversation.contact_id is None or conversation.connection_id is None:
|
||
raise CallConflict(t("calls.no_client_identity"))
|
||
identities = list(
|
||
ConnectionIdentity.objects.filter(
|
||
contact_id=conversation.contact_id,
|
||
connection_id=conversation.connection_id,
|
||
)
|
||
.order_by("id")[:2]
|
||
)
|
||
if len(identities) != 1:
|
||
raise CallConflict(t("calls.identity_missing_or_ambiguous"))
|
||
return identities[0]
|
||
|
||
|
||
def _check_call_creation_conflicts(*, conversation: Conversation, initiator) -> None:
|
||
if conversation.lifecycle != LifecycleState.OPEN:
|
||
raise CallConflict(t("calls.only_in_open_conversation"))
|
||
if CallSession.objects.filter(
|
||
conversation=conversation,
|
||
status__in=UNFINISHED_CALL_STATUSES,
|
||
).exists():
|
||
raise CallConflict(t("calls.unfinished_call_exists"))
|
||
if CallSession.objects.filter(
|
||
initiated_by=initiator,
|
||
status__in=UNFINISHED_CALL_STATUSES,
|
||
).exists():
|
||
raise CallConflict(t("calls.operator_busy"))
|
||
if (
|
||
conversation.control_mode == ControlMode.HUMAN
|
||
and conversation.assigned_operator_id
|
||
and conversation.assigned_operator_id != initiator.id
|
||
):
|
||
raise CallConflict(t("calls.conversation_taken"))
|
||
|
||
|
||
@transaction.atomic
|
||
def create_call_request(
|
||
*, context: TenantContext, conversation_id: int, kind: str = CallKind.AUDIO
|
||
) -> CreatedCall:
|
||
initiator = context.actor_user
|
||
if initiator is None or context.membership is None:
|
||
raise CallAccessDenied(t("calls.operator_context_required"))
|
||
conversation = (
|
||
Conversation.objects.select_for_update()
|
||
.select_related("channel")
|
||
.get(id=conversation_id, organization=context.organization)
|
||
)
|
||
ensure_conversation_call_access(user=context.membership, conversation=conversation)
|
||
if not call_allowed(conversation.connection, kind):
|
||
raise CallAccessDenied(
|
||
t("calls.video_off_entry_point") if kind == CallKind.VIDEO else t("calls.calls_off_entry_point")
|
||
)
|
||
identity = _conversation_identity(conversation)
|
||
_check_call_creation_conflicts(conversation=conversation, initiator=initiator)
|
||
|
||
if conversation.control_mode != ControlMode.HUMAN or conversation.assigned_operator_id != initiator.id:
|
||
try:
|
||
claim_locked_conversation(context=context, conversation=conversation)
|
||
except ClaimError as error:
|
||
raise CallConflict(str(error)) from error
|
||
|
||
invite_token, token_hash = issue_invite_token()
|
||
try:
|
||
with transaction.atomic():
|
||
call = CallSession.objects.create(
|
||
organization_id=conversation.organization_id,
|
||
conversation=conversation,
|
||
initiated_by=initiator,
|
||
delivery_connection_id=conversation.connection_id,
|
||
kind=kind,
|
||
)
|
||
except IntegrityError as error:
|
||
raise CallConflict(t("calls.second_call_rejected")) from error
|
||
|
||
CallInvite.objects.create(
|
||
organization=context.organization,
|
||
call_session=call,
|
||
connection_identity=identity,
|
||
token_hash=token_hash,
|
||
expires_at=timezone.now() + timedelta(seconds=settings.CHATBALLS_CALL_INVITE_TTL_SECONDS),
|
||
)
|
||
CallParticipant.objects.bulk_create(
|
||
[
|
||
CallParticipant(
|
||
organization=context.organization,
|
||
call_session=call,
|
||
side=ParticipantSide.STAFF,
|
||
user=initiator,
|
||
),
|
||
CallParticipant(
|
||
organization=context.organization,
|
||
call_session=call,
|
||
side=ParticipantSide.CUSTOMER,
|
||
connection_identity=identity,
|
||
),
|
||
]
|
||
)
|
||
initiator_label = getattr(initiator, "full_name", "") or initiator.email
|
||
call_word = "аудиозвонок" if call.kind == CallKind.AUDIO else "видеозвонок"
|
||
Message.objects.create(
|
||
conversation=conversation,
|
||
author_type=MessageAuthor.SYSTEM,
|
||
system_event=SystemEvent.CALL_REQUESTED,
|
||
# Вид звонка — параметр, а не часть кода: фразу собирает интерфейс, и
|
||
# «аудио» против «видео» там отдельным словом словаря.
|
||
system_params={"operator": initiator_label, "kind": call.kind},
|
||
text=f"Оператор {initiator_label} запросил {call_word}",
|
||
)
|
||
if conversation.connection.provider == IntegrationProvider.WEB:
|
||
# Web Chat: приглашение забирает виджет поллингом, внешней отправки нет.
|
||
call.invite.delivery_status = InviteDeliveryStatus.SENT
|
||
call.invite.save(update_fields=["delivery_status"])
|
||
call = transition_call(call_session_id=call.id, target_status=CallStatus.RINGING)
|
||
else:
|
||
# TG/MAX: кнопка со ссылкой /calls/<token> уходит через outbox с ретраями.
|
||
enqueue_event(
|
||
DomainEvent(
|
||
aggregate_type="call_session",
|
||
aggregate_id=str(call.id),
|
||
event_type=CALL_INVITE_SEND,
|
||
payload={"callSessionId": str(call.id)},
|
||
tenant_context=context,
|
||
)
|
||
)
|
||
staff_token = issue_call_access_token(
|
||
call_session_id=call.id,
|
||
side=ParticipantSide.STAFF,
|
||
subject_id=str(initiator.id),
|
||
)
|
||
return CreatedCall(call_session=call, invite_token=invite_token, staff_access_token=staff_token)
|
||
|
||
|
||
def issue_staff_access_token(*, context: TenantContext, call_session: CallSession) -> str:
|
||
user = context.actor_user
|
||
if user is None or context.membership is None:
|
||
raise CallAccessDenied(t("calls.operator_context_required"))
|
||
ensure_call_access(user=context.membership, call_session=call_session)
|
||
participant_exists = call_session.participants.filter(side=ParticipantSide.STAFF, user=user).exists()
|
||
if not participant_exists:
|
||
raise CallConflict(t("calls.not_a_participant"))
|
||
if call_session.status in TERMINAL_CALL_STATUSES:
|
||
raise CallConflict(t("calls.already_ended"))
|
||
return issue_call_access_token(
|
||
call_session_id=call_session.id,
|
||
side=ParticipantSide.STAFF,
|
||
subject_id=str(user.id),
|
||
)
|
||
|
||
|
||
def cancel_call(*, context: TenantContext, call_session: CallSession) -> CallSession:
|
||
"""Оператор закончил звонок — из любой фазы, в которой тот ещё жив.
|
||
|
||
Кнопка у оператора одна и означает «прекратить»: до ответа клиента это
|
||
отмена, после — завершение. Раньше здесь был только переход в CANCELLED, и
|
||
он разрешён лишь из REQUESTED/RINGING: клиент принял звонок, соединение не
|
||
установилось (частый случай за NAT), оператор жмёт «завершить» — и получает
|
||
409, а звонок остаётся висеть. Фазу выбирает `finish_call`, тот же код, что
|
||
и у клиента.
|
||
"""
|
||
user = context.actor_user
|
||
if user is None or context.membership is None:
|
||
raise CallAccessDenied(t("calls.operator_context_required"))
|
||
ensure_call_access(user=context.membership, call_session=call_session)
|
||
if not call_session.participants.filter(side=ParticipantSide.STAFF, user=user).exists():
|
||
raise CallConflict(t("calls.not_a_participant"))
|
||
try:
|
||
# Повторный вызов идемпотентен: у завершённого звонка finish_call
|
||
# возвращает его как есть.
|
||
return finish_call(call_session_id=call_session.id, side=ParticipantSide.STAFF)
|
||
except CallInvalidTransition as error:
|
||
raise CallConflict(t("calls.cannot_cancel")) from error
|
||
|
||
|
||
def active_call_for_conversation(conversation: Conversation) -> CallSession | None:
|
||
return (
|
||
CallSession.objects.filter(
|
||
conversation=conversation,
|
||
status__in=UNFINISHED_CALL_STATUSES,
|
||
)
|
||
.select_related("initiated_by")
|
||
.prefetch_related("participants")
|
||
.first()
|
||
)
|
||
|
||
|
||
# --- Web Chat: приглашение доставляется поллингом виджета по session identity ---
|
||
|
||
|
||
def webchat_active_call(identity: ConnectionIdentity) -> CallSession | None:
|
||
return (
|
||
CallSession.objects.filter(
|
||
invite__connection_identity=identity,
|
||
status__in=UNFINISHED_CALL_STATUSES,
|
||
)
|
||
.select_related("invite", "initiated_by")
|
||
.first()
|
||
)
|
||
|
||
|
||
@transaction.atomic
|
||
def open_call_for_identity(*, identity: ConnectionIdentity) -> ResolvedInvite:
|
||
call = webchat_active_call(identity)
|
||
if call is None:
|
||
raise CallTokenError(t("webchat.no_active_invite"))
|
||
invite = CallInvite.objects.select_for_update().get(call_session=call)
|
||
if invite.expires_at <= timezone.now():
|
||
raise CallTokenError(t("webchat.no_active_invite"))
|
||
if invite.opened_at is None:
|
||
invite.opened_at = timezone.now()
|
||
invite.save(update_fields=["opened_at"])
|
||
access_token = issue_call_access_token(
|
||
call_session_id=call.id,
|
||
side=ParticipantSide.CUSTOMER,
|
||
subject_id=str(invite.id),
|
||
)
|
||
return ResolvedInvite(invite=invite, customer_access_token=access_token)
|
||
|
||
|
||
def decline_call_for_identity(*, identity: ConnectionIdentity) -> CallSession:
|
||
call = webchat_active_call(identity)
|
||
if call is None:
|
||
raise CallTokenError(t("webchat.no_active_invite"))
|
||
try:
|
||
return transition_call(
|
||
call_session_id=call.id,
|
||
target_status=CallStatus.DECLINED,
|
||
ended_by=CallEndedBy.CUSTOMER,
|
||
)
|
||
except CallInvalidTransition as error:
|
||
raise CallConflict(t("calls.invite_cannot_decline")) from error
|