From 3659ce564a13981d151b81beb4617dee6752dd1b Mon Sep 17 00:00:00 2001 From: Andrey Date: Tue, 8 Sep 2026 13:58:39 +0300 Subject: [PATCH] =?UTF-8?q?:lock:=20fix(calls):=20=D1=81=D0=BE=D0=BA=D0=B5?= =?UTF-8?q?=D1=82=20=D1=81=D0=B8=D0=B3=D0=BD=D0=B0=D0=BB=D0=B8=D0=BD=D0=B3?= =?UTF-8?q?=D0=B0=20=D0=BD=D0=B5=20=D0=B6=D0=B4=D1=91=D1=82=20=D0=B0=D1=83?= =?UTF-8?q?=D1=82=D0=B5=D0=BD=D1=82=D0=B8=D1=84=D0=B8=D0=BA=D0=B0=D1=86?= =?UTF-8?q?=D0=B8=D0=B8=20=D0=B2=D0=B5=D1=87=D0=BD=D0=BE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Константа AUTH_TIMEOUT_CLOSE намекала на срок, но кода не было: соединение принималось до аутентификации (иначе клиенту некуда прислать токен) и дальше ждало первое сообщение сколько угодно долго. Неаутентифицированный клиент так держал сокет и запись в channel layer. Срок — 10 секунд, снимается при успешной аутентификации и при отключении. Два теста. Co-Authored-By: Claude Opus 5 --- apps/backend/chatballs/calls/consumers.py | 25 +++++++++ .../chatballs/calls/tests/test_signaling.py | 52 ++++++++++++++++++- 2 files changed, 76 insertions(+), 1 deletion(-) diff --git a/apps/backend/chatballs/calls/consumers.py b/apps/backend/chatballs/calls/consumers.py index bdbe2fe..60224eb 100644 --- a/apps/backend/chatballs/calls/consumers.py +++ b/apps/backend/chatballs/calls/consumers.py @@ -9,6 +9,7 @@ - reconnect WebSocket не создаёт новую CallSession — только presence-статусы. """ +import asyncio import logging from channels.db import database_sync_to_async @@ -25,6 +26,11 @@ from chatballs.tenancy.database import run_tenant_operation logger = logging.getLogger(__name__) AUTH_TIMEOUT_CLOSE = 4401 +# Сколько ждём первое сообщение {"type": "auth"}. Соединение принимается до +# аутентификации (иначе клиенту некуда прислать токен), поэтому без срока +# неаутентифицированный клиент держал бы сокет и запись в channel layer +# сколько угодно долго. +AUTH_TIMEOUT_SECONDS = 10 # События, которые сервер только ретранслирует второму участнику. RELAY_TYPES = {"webrtc.offer", "webrtc.answer", "webrtc.ice_candidate", "participant.media_state"} SEEN_COMMANDS_LIMIT = 512 @@ -38,6 +44,21 @@ class CallSignalingConsumer(AsyncJsonWebsocketConsumer): self.group = None self._seen_commands: set[str] = set() await self.accept() + self._auth_deadline = asyncio.create_task(self._close_unless_authenticated()) + + async def _close_unless_authenticated(self) -> None: + try: + await asyncio.sleep(AUTH_TIMEOUT_SECONDS) + except asyncio.CancelledError: + return + if self.call_id is None: + await self.close(code=AUTH_TIMEOUT_CLOSE) + + def _cancel_auth_deadline(self) -> None: + deadline = getattr(self, "_auth_deadline", None) + if deadline is not None and not deadline.done(): + deadline.cancel() + self._auth_deadline = None async def receive_json(self, content: dict, **kwargs) -> None: msg_type = content.get("type") @@ -81,6 +102,7 @@ class CallSignalingConsumer(AsyncJsonWebsocketConsumer): # незнакомые типы игнорируются без разрыва соединения async def disconnect(self, code: int) -> None: + self._cancel_auth_deadline() if self.group is None: return await self.channel_layer.group_discard(self.group, self.channel_name) @@ -94,6 +116,7 @@ class CallSignalingConsumer(AsyncJsonWebsocketConsumer): async def _authenticate(self, msg_type, content: dict) -> None: if msg_type != "auth": + self._cancel_auth_deadline() await self.close(code=AUTH_TIMEOUT_CLOSE) return token = str(content.get("token", "")) @@ -102,9 +125,11 @@ class CallSignalingConsumer(AsyncJsonWebsocketConsumer): lambda: authorize_call_access_context(token=token) )() except CallTokenError: + self._cancel_auth_deadline() await self.send_json({"type": "error", "code": "AUTH_FAILED"}) await self.close(code=AUTH_TIMEOUT_CLOSE) return + self._cancel_auth_deadline() self.call_id = call.id self.side = claims.side self.tenant_context = context diff --git a/apps/backend/chatballs/calls/tests/test_signaling.py b/apps/backend/chatballs/calls/tests/test_signaling.py index 620283d..001ce5c 100644 --- a/apps/backend/chatballs/calls/tests/test_signaling.py +++ b/apps/backend/chatballs/calls/tests/test_signaling.py @@ -2,14 +2,16 @@ relay только между участниками звонка, переходы CONNECTING/ACTIVE/ENDED, reconnect без новой CallSession, поздние события игнорируются.""" +import asyncio from datetime import timedelta +from unittest import mock from asgiref.sync import async_to_sync from channels.testing import WebsocketCommunicator from django.test import TransactionTestCase from django.utils import timezone -from chatballs_backend.asgi import application +from chatballs.calls import consumers from chatballs.calls.models import ( CallParticipant, CallSession, @@ -23,6 +25,7 @@ from chatballs.calls.services import ( ) from chatballs.calls.tests.helpers import CallDomainMixin, create_call_request, expire_stale_calls from chatballs.conversations.models import Message +from chatballs_backend.asgi import application WS_PATH = "/ws/calls/" @@ -245,3 +248,50 @@ class ReconnectSweepTests(SignalingTestCase): expire_stale_calls(self.organization) call = CallSession.objects.get() self.assertEqual(call.status, CallStatus.ACTIVE) + + +class AuthDeadlineTests(CallDomainMixin, TransactionTestCase): + """Соединение принимается до аутентификации — значит, ждать оно должно не вечно.""" + + def test_unauthenticated_socket_is_closed_on_deadline(self) -> None: + async def scenario(): + communicator = WebsocketCommunicator(application, WS_PATH) + connected, _ = await communicator.connect() + assert connected + # Токен не присылаем вовсе: сокет обязан закрыться сам. + output = await communicator.receive_output(timeout=5) + await communicator.disconnect() + return output + + with mock.patch.object(consumers, "AUTH_TIMEOUT_SECONDS", 0.1): + output = async_to_sync(scenario)() + + self.assertEqual(output["type"], "websocket.close") + self.assertEqual(output["code"], consumers.AUTH_TIMEOUT_CLOSE) + + def test_authenticated_socket_survives_the_deadline(self) -> None: + created, customer_token = None, None + + def prepare(): + create_call_request(conversation_id=self.conversation.id, initiator=self.owner) + return open_call_for_identity(identity=self.identity).customer_access_token + + customer_token = prepare() + + async def scenario(token: str): + communicator = WebsocketCommunicator(application, WS_PATH) + connected, _ = await communicator.connect() + assert connected + await communicator.send_json_to({"type": "auth", "token": token}) + state = await communicator.receive_json_from() + # Пережидаем срок: у аутентифицированного соединения он снят. + await asyncio.sleep(0.3) + nothing_left = await communicator.receive_nothing(timeout=0.2) + await communicator.disconnect() + return state, nothing_left + + with mock.patch.object(consumers, "AUTH_TIMEOUT_SECONDS", 0.1): + state, nothing_left = async_to_sync(scenario)(customer_token) + + self.assertEqual(state["type"], "call.state") + self.assertTrue(nothing_left)