mirror of
https://github.com/dartdavros/chatballs.git
synced 2026-10-05 01:14:58 +03:00
🔒 fix(calls): сокет сигналинга не ждёт аутентификации вечно
Константа AUTH_TIMEOUT_CLOSE намекала на срок, но кода не было: соединение принималось до аутентификации (иначе клиенту некуда прислать токен) и дальше ждало первое сообщение сколько угодно долго. Неаутентифицированный клиент так держал сокет и запись в channel layer. Срок — 10 секунд, снимается при успешной аутентификации и при отключении. Два теста. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
1 parent
037be19786
commit
3659ce564a
2 files changed
+76
-1
No files matched your search
@@ -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
|
||||
|
||||
@@ -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)
|
||||
Reference in new issue
Block a user