mirror of
https://github.com/dartdavros/chatballs.git
synced 2026-10-05 09:14:58 +03:00
⚡ perf(realtime): обновление по событию вместо опроса
Клиент опрашивал сервер: список раз в четыре секунды, карточку и дельту истории — раз в три; на оператора выходило под сорок запросов в минуту, и новое сообщение всё равно появлялось с задержкой. - WebSocket-канал организации (channels уже стоял ради сигналинга звонков): аутентификация — сессией того же SPA, организация — в адресе, как в HTTP; - событие несёт только повод обновиться, данные клиент забирает обычным запросом: проверка видимости остаётся в одном месте, и канал не может в ней ошибиться. Событие инбокса не содержит идентификаторов — сотрудник видит не все диалоги организации; на события диалога подписка отдельная, и сервер проверяет видимость перед ней; - события шлют те же сигналы, что держат свежесть диалога: одно место на семь мест создания сообщений и на все изменения состояния; - опрос остался запасным путём и замедляется до тридцати секунд, пока канал жив: при обрыве всё возвращается к прежнему поведению само. Проверено сквозь шлюз: апгрейд проходит до консьюмера, анонимное соединение отклоняется кодом 403. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
1 parent
f3cbf1658c
commit
03825dfbd5
11 files changed
+536
-18
No files matched your search
@@ -0,0 +1,116 @@
|
||||
"""WebSocket оповещений о диалогах (см. chatballs.conversations.realtime).
|
||||
|
||||
Правила:
|
||||
- аутентификация — сессией того же SPA (AuthMiddlewareStack), отдельного токена
|
||||
нет: сокет открывает тот же браузер, что и REST;
|
||||
- организация берётся из адреса, как и в HTTP-слое, и обязана совпадать с
|
||||
членством пользователя;
|
||||
- подписка на диалог возможна только после проверки видимости — той же, что у
|
||||
REST (ADR-CHATBALLS-0043);
|
||||
- события не содержат содержимого: клиент по ним перезапрашивает данные и
|
||||
получает ровно то, что ему позволено.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import uuid
|
||||
|
||||
from channels.db import database_sync_to_async
|
||||
from channels.generic.websocket import AsyncJsonWebsocketConsumer
|
||||
|
||||
from chatballs.conversations.models import Conversation
|
||||
from chatballs.conversations.realtime import conversation_group, inbox_group
|
||||
from chatballs.conversations.selectors import conversation_is_visible
|
||||
from chatballs.identity.models import Organization, OrganizationMembership
|
||||
from chatballs.tenancy.database import tenant_atomic
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
NOT_A_MEMBER_CLOSE = 4403
|
||||
|
||||
|
||||
class ConversationEventsConsumer(AsyncJsonWebsocketConsumer):
|
||||
async def connect(self) -> None:
|
||||
self.organization_id: int | None = None
|
||||
self.membership_id: int | None = None
|
||||
self.watched: str | None = None
|
||||
user = self.scope.get("user")
|
||||
if user is None or not user.is_authenticated:
|
||||
await self.close(code=NOT_A_MEMBER_CLOSE)
|
||||
return
|
||||
raw_public_id = self.scope["url_route"]["kwargs"]["organization_public_id"]
|
||||
resolved = await self._membership(user.id, raw_public_id)
|
||||
if resolved is None:
|
||||
await self.close(code=NOT_A_MEMBER_CLOSE)
|
||||
return
|
||||
self.organization_id, self.membership_id = resolved
|
||||
await self.channel_layer.group_add(inbox_group(self.organization_id), self.channel_name)
|
||||
await self.accept()
|
||||
|
||||
async def disconnect(self, code: int) -> None:
|
||||
if self.organization_id is not None:
|
||||
await self.channel_layer.group_discard(
|
||||
inbox_group(self.organization_id), self.channel_name
|
||||
)
|
||||
if self.watched is not None:
|
||||
await self.channel_layer.group_discard(self.watched, self.channel_name)
|
||||
|
||||
async def receive_json(self, content: dict, **kwargs) -> None:
|
||||
"""Клиент сообщает, какой диалог открыт: событий по нему он и ждёт."""
|
||||
if content.get("type") != "watch" or self.organization_id is None:
|
||||
return
|
||||
conversation_id = content.get("conversationId")
|
||||
if self.watched is not None:
|
||||
await self.channel_layer.group_discard(self.watched, self.channel_name)
|
||||
self.watched = None
|
||||
if not isinstance(conversation_id, int):
|
||||
return
|
||||
if not await self._may_watch(conversation_id):
|
||||
# Молча: отсутствие подписки и отсутствие прав снаружи неразличимы.
|
||||
return
|
||||
self.watched = conversation_group(conversation_id)
|
||||
await self.channel_layer.group_add(self.watched, self.channel_name)
|
||||
# Подтверждение — не вежливость: пока его нет, события диалога ещё могут
|
||||
# пройти мимо, и клиенту после переподключения нужно знать, с какого
|
||||
# момента лента снова живая.
|
||||
await self.send_json({"type": "watching", "conversationId": conversation_id})
|
||||
|
||||
async def fanout(self, event: dict) -> None:
|
||||
await self.send_json(event["payload"])
|
||||
|
||||
@database_sync_to_async
|
||||
def _membership(self, user_id: int, raw_public_id: str) -> tuple[int, int] | None:
|
||||
try:
|
||||
organization = Organization.objects.get(public_id=uuid.UUID(str(raw_public_id)))
|
||||
except (ValueError, Organization.DoesNotExist):
|
||||
return None
|
||||
with tenant_atomic(organization.pk):
|
||||
membership = (
|
||||
OrganizationMembership.objects.select_related("user")
|
||||
.filter(
|
||||
organization_id=organization.pk,
|
||||
user_id=user_id,
|
||||
blocked_at__isnull=True,
|
||||
user__is_active=True,
|
||||
)
|
||||
.first()
|
||||
)
|
||||
if membership is None:
|
||||
return None
|
||||
return organization.pk, membership.pk
|
||||
|
||||
@database_sync_to_async
|
||||
def _may_watch(self, conversation_id: int) -> bool:
|
||||
with tenant_atomic(self.organization_id):
|
||||
membership = (
|
||||
OrganizationMembership.objects.select_related("user")
|
||||
.filter(pk=self.membership_id)
|
||||
.first()
|
||||
)
|
||||
conversation = Conversation.objects.filter(
|
||||
id=conversation_id, organization_id=self.organization_id
|
||||
).first()
|
||||
if membership is None or conversation is None:
|
||||
return False
|
||||
return conversation_is_visible(actor=membership, conversation=conversation)
|
||||
@@ -0,0 +1,59 @@
|
||||
"""Оповещения о диалогах: сервер сообщает, что изменилось, а не что показать.
|
||||
|
||||
Клиент до сих пор опрашивал сервер: список раз в четыре секунды, карточку и
|
||||
дельту истории — раз в три. Окна сделали каждый запрос дешёвым, но не
|
||||
бесплатным: на оператора выходило под сорок запросов в минуту, и новое
|
||||
сообщение всё равно появлялось с задержкой.
|
||||
|
||||
Событие несёт только повод обновиться. Данные клиент забирает обычным REST'ом,
|
||||
где и живёт проверка видимости: канал не решает, кому что показывать, и потому
|
||||
не может ошибиться в этом. По той же причине событие инбокса не содержит
|
||||
идентификаторов — сотрудник видит не все диалоги организации.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
|
||||
from asgiref.sync import async_to_sync
|
||||
from channels.layers import get_channel_layer
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
INBOX_EVENT = "inbox.changed"
|
||||
CONVERSATION_EVENT = "conversation.changed"
|
||||
|
||||
|
||||
def inbox_group(organization_id: int) -> str:
|
||||
return f"inbox.{organization_id}"
|
||||
|
||||
|
||||
def conversation_group(conversation_id: int) -> str:
|
||||
return f"conv.{conversation_id}"
|
||||
|
||||
|
||||
def _publish(group: str, payload: dict[str, object]) -> None:
|
||||
"""Оповещение — вспомогательный путь: его сбой не должен ронять запись.
|
||||
|
||||
Сообщение уже сохранено к моменту отправки; если канал недоступен, клиент
|
||||
узнает об изменении следующим опросом — он остаётся как запасной путь.
|
||||
"""
|
||||
layer = get_channel_layer()
|
||||
if layer is None:
|
||||
return
|
||||
try:
|
||||
async_to_sync(layer.group_send)(group, {"type": "fanout", "payload": payload})
|
||||
except Exception: # noqa: BLE001 — канал не должен ломать сохранение
|
||||
logger.warning("realtime fanout failed for %s", group, exc_info=True)
|
||||
|
||||
|
||||
def notify_inbox_changed(organization_id: int) -> None:
|
||||
_publish(inbox_group(organization_id), {"type": INBOX_EVENT})
|
||||
|
||||
|
||||
def notify_conversation_changed(conversation_id: int, *, organization_id: int) -> None:
|
||||
_publish(
|
||||
conversation_group(conversation_id),
|
||||
{"type": CONVERSATION_EVENT, "conversationId": conversation_id},
|
||||
)
|
||||
notify_inbox_changed(organization_id)
|
||||
@@ -0,0 +1,14 @@
|
||||
from channels.auth import AuthMiddlewareStack
|
||||
from django.urls import path
|
||||
|
||||
from chatballs.conversations.consumers import ConversationEventsConsumer
|
||||
|
||||
# Оповещения о диалогах аутентифицируются сессией того же SPA — отдельного
|
||||
# токена, как у сигналинга звонков, здесь не нужно: сокет открывает тот же
|
||||
# браузер. Организация стоит в адресе, как и во всём HTTP-слое.
|
||||
websocket_urlpatterns = [
|
||||
path(
|
||||
"ws/organizations/<uuid:organization_public_id>/conversations/",
|
||||
AuthMiddlewareStack(ConversationEventsConsumer.as_asgi()),
|
||||
),
|
||||
]
|
||||
@@ -1,16 +1,25 @@
|
||||
"""Сигналы домена диалогов.
|
||||
|
||||
Единственная задача — держать `Conversation.last_message_at` в согласии с
|
||||
лентой. Сообщения создаются в семи местах (приём из каналов, ответ оператора,
|
||||
голосовые, файлы, события звонков, демо-данные), поэтому обновление живёт в
|
||||
сигнале: иначе достаточно одного забытого места, чтобы диалог перестал
|
||||
подниматься в инбоксе.
|
||||
Две задачи, и обе — про «после записи»:
|
||||
|
||||
* держать `Conversation.last_message_at` в согласии с лентой;
|
||||
* оповещать открытые интерфейсы, что диалог изменился.
|
||||
|
||||
Сообщения создаются в семи местах (приём из каналов, ответ оператора,
|
||||
голосовые, файлы, события звонков, демо-данные), а состояние диалога меняют
|
||||
ещё и перехват, возврат AI, метки, приоритет, архив. Поэтому и то и другое
|
||||
живёт в сигналах: иначе достаточно одного забытого места, чтобы диалог перестал
|
||||
подниматься в инбоксе или чтобы у оператора не обновился экран.
|
||||
"""
|
||||
|
||||
from django.db.models.signals import post_save
|
||||
from django.dispatch import receiver
|
||||
|
||||
from chatballs.conversations.models import Conversation, Message
|
||||
from chatballs.conversations.realtime import (
|
||||
notify_conversation_changed,
|
||||
notify_inbox_changed,
|
||||
)
|
||||
|
||||
|
||||
@receiver(post_save, sender=Message, dispatch_uid="conversations.touch_last_message_at")
|
||||
@@ -25,3 +34,21 @@ def touch_last_message_at(sender, instance: Message, created: bool, **kwargs) ->
|
||||
Conversation.objects.filter(
|
||||
id=instance.conversation_id, last_message_at__lt=instance.created_at
|
||||
).update(last_message_at=instance.created_at)
|
||||
|
||||
|
||||
@receiver(post_save, sender=Message, dispatch_uid="conversations.notify_message")
|
||||
def notify_message(sender, instance: Message, created: bool, **kwargs) -> None:
|
||||
if not created:
|
||||
return
|
||||
notify_conversation_changed(
|
||||
instance.conversation_id, organization_id=instance.organization_id
|
||||
)
|
||||
|
||||
|
||||
@receiver(post_save, sender=Conversation, dispatch_uid="conversations.notify_conversation")
|
||||
def notify_conversation(sender, instance: Conversation, created: bool, **kwargs) -> None:
|
||||
"""Перехват, возврат AI, метки, приоритет, архив — всё это запись строки."""
|
||||
if created:
|
||||
notify_inbox_changed(instance.organization_id)
|
||||
return
|
||||
notify_conversation_changed(instance.id, organization_id=instance.organization_id)
|
||||
@@ -0,0 +1,168 @@
|
||||
"""Оповещения о диалогах: кто их получает и что в них лежит.
|
||||
|
||||
Канал не решает, кому что показывать: событие несёт только повод обновиться, а
|
||||
данные клиент забирает REST'ом, где и живёт проверка видимости. Поэтому здесь
|
||||
проверяется ровно две вещи — что событие доходит до того, кто вправе его ждать,
|
||||
и что подписаться на чужой диалог нельзя.
|
||||
"""
|
||||
|
||||
from channels.testing import WebsocketCommunicator
|
||||
from django.test import TransactionTestCase
|
||||
|
||||
from chatballs.channels.models import Channel
|
||||
from chatballs.conversations.models import (
|
||||
Contact,
|
||||
Conversation,
|
||||
LifecycleState,
|
||||
Message,
|
||||
MessageAuthor,
|
||||
)
|
||||
from chatballs.identity.bootstrap import bootstrap_owner
|
||||
from chatballs.identity.group_models import EmployeeGroup
|
||||
from chatballs.identity.models import (
|
||||
EmployeeRole,
|
||||
HumanUser,
|
||||
Organization,
|
||||
OrganizationMembership,
|
||||
)
|
||||
from chatballs_backend.asgi_app import application
|
||||
|
||||
|
||||
class ConversationEventsTests(TransactionTestCase):
|
||||
def setUp(self) -> None:
|
||||
result = bootstrap_owner(email="owner@example.com", password="temporary-password")
|
||||
self.owner = result.owner
|
||||
self.organization = Organization.objects.get(slug="demo")
|
||||
self.channel = Channel.objects.create(
|
||||
organization=self.organization, code="line", name="Линия"
|
||||
)
|
||||
contact = Contact.objects.create(organization=self.organization, name="Иван")
|
||||
self.conversation = Conversation.objects.create(
|
||||
organization=self.organization,
|
||||
channel=self.channel,
|
||||
contact=contact,
|
||||
lifecycle=LifecycleState.OPEN,
|
||||
)
|
||||
|
||||
def _url(self) -> str:
|
||||
return f"/ws/organizations/{self.organization.public_id}/conversations/"
|
||||
|
||||
def _communicator(self, user) -> WebsocketCommunicator:
|
||||
communicator = WebsocketCommunicator(application, self._url())
|
||||
communicator.scope["user"] = user
|
||||
return communicator
|
||||
|
||||
async def _connect(self, user) -> WebsocketCommunicator:
|
||||
communicator = self._communicator(user)
|
||||
connected, _ = await communicator.connect()
|
||||
self.assertTrue(connected)
|
||||
return communicator
|
||||
|
||||
async def test_stranger_is_not_connected(self) -> None:
|
||||
outsider = await self._create_outsider()
|
||||
communicator = self._communicator(outsider)
|
||||
connected, code = await communicator.connect()
|
||||
self.assertFalse(connected)
|
||||
self.assertEqual(code, 4403)
|
||||
|
||||
async def test_new_message_reaches_the_open_dialog(self) -> None:
|
||||
communicator = await self._connect(self.owner)
|
||||
await communicator.send_json_to(
|
||||
{"type": "watch", "conversationId": self.conversation.id}
|
||||
)
|
||||
# Подписка подтверждается: до подтверждения события диалога могли бы
|
||||
# пройти мимо.
|
||||
self.assertEqual(
|
||||
await communicator.receive_json_from(timeout=5),
|
||||
{"type": "watching", "conversationId": self.conversation.id},
|
||||
)
|
||||
await self._post_message()
|
||||
# Подписчик получает и событие диалога, и общее событие инбокса; порядок
|
||||
# между группами не определён, поэтому проверяется состав.
|
||||
events = [await communicator.receive_json_from(timeout=5) for _ in range(2)]
|
||||
changed = next(event for event in events if event["type"] == "conversation.changed")
|
||||
self.assertEqual(changed["conversationId"], self.conversation.id)
|
||||
await communicator.disconnect()
|
||||
|
||||
async def test_inbox_event_carries_no_identifiers(self) -> None:
|
||||
# Сотрудник видит не все диалоги организации, поэтому в событии инбокса
|
||||
# не должно быть ни идентификаторов, ни содержимого.
|
||||
communicator = await self._connect(self.owner)
|
||||
await self._post_message()
|
||||
# Без подписки на диалог доходит только событие инбокса — и в нём нет
|
||||
# ничего, кроме самого повода обновиться.
|
||||
event = await communicator.receive_json_from(timeout=5)
|
||||
self.assertEqual(event, {"type": "inbox.changed"})
|
||||
self.assertTrue(await communicator.receive_nothing(timeout=1))
|
||||
await communicator.disconnect()
|
||||
|
||||
async def test_watch_is_refused_for_a_dialog_outside_visibility(self) -> None:
|
||||
employee, hidden = await self._create_employee_and_hidden_dialog()
|
||||
communicator = await self._connect(employee)
|
||||
await communicator.send_json_to({"type": "watch", "conversationId": hidden.id})
|
||||
await self._post_message(conversation=hidden)
|
||||
# На чужой диалог подписки нет: до клиента доходит только событие
|
||||
# инбокса, без идентификатора.
|
||||
event = await communicator.receive_json_from(timeout=5)
|
||||
self.assertEqual(event, {"type": "inbox.changed"})
|
||||
self.assertTrue(await communicator.receive_nothing(timeout=1))
|
||||
await communicator.disconnect()
|
||||
|
||||
async def _post_message(self, conversation: Conversation | None = None) -> None:
|
||||
from channels.db import database_sync_to_async
|
||||
|
||||
target = conversation or self.conversation
|
||||
|
||||
@database_sync_to_async
|
||||
def create() -> None:
|
||||
Message.objects.create(
|
||||
conversation=target, author_type=MessageAuthor.CONTACT, text="Здравствуйте"
|
||||
)
|
||||
|
||||
await create()
|
||||
|
||||
async def _create_outsider(self):
|
||||
from channels.db import database_sync_to_async
|
||||
|
||||
@database_sync_to_async
|
||||
def create():
|
||||
other = Organization.objects.create(name="Другая", slug="other-org")
|
||||
user = HumanUser.objects.create_user(
|
||||
email="stranger@example.com", password="Password-123"
|
||||
)
|
||||
OrganizationMembership.objects.create(
|
||||
user=user, organization=other, role=EmployeeRole.OWNER, position_title="Owner"
|
||||
)
|
||||
return user
|
||||
|
||||
return await create()
|
||||
|
||||
async def _create_employee_and_hidden_dialog(self):
|
||||
from channels.db import database_sync_to_async
|
||||
|
||||
@database_sync_to_async
|
||||
def create():
|
||||
user = HumanUser.objects.create_user(
|
||||
email="operator@example.com", password="Password-123"
|
||||
)
|
||||
OrganizationMembership.objects.create(
|
||||
user=user,
|
||||
organization=self.organization,
|
||||
role=EmployeeRole.EMPLOYEE,
|
||||
position_title="Оператор",
|
||||
)
|
||||
# Диалог чужой группы: сотрудник в неё не входит и видеть его не должен.
|
||||
group = EmployeeGroup.objects.create(
|
||||
organization=self.organization, name="Закрытая группа"
|
||||
)
|
||||
contact = Contact.objects.create(organization=self.organization, name="Пётр")
|
||||
hidden = Conversation.objects.create(
|
||||
organization=self.organization,
|
||||
channel=self.channel,
|
||||
contact=contact,
|
||||
group=group,
|
||||
lifecycle=LifecycleState.OPEN,
|
||||
)
|
||||
return user, hidden
|
||||
|
||||
return await create()
|
||||
@@ -9,11 +9,14 @@ django_asgi_app = get_asgi_application()
|
||||
|
||||
from channels.routing import ProtocolTypeRouter, URLRouter # noqa: E402
|
||||
|
||||
from chatballs.calls.routing import websocket_urlpatterns # noqa: E402
|
||||
from chatballs.calls.routing import websocket_urlpatterns as call_routes # noqa: E402
|
||||
from chatballs.conversations.routing import ( # noqa: E402
|
||||
websocket_urlpatterns as conversation_routes,
|
||||
)
|
||||
|
||||
application = ProtocolTypeRouter(
|
||||
{
|
||||
"http": django_asgi_app,
|
||||
"websocket": URLRouter(websocket_urlpatterns),
|
||||
"websocket": URLRouter(call_routes + conversation_routes),
|
||||
}
|
||||
)
|
||||
@@ -69,6 +69,17 @@ export function resolveApiUrl(path: string): string {
|
||||
return `${configuredApiBase}${path}`;
|
||||
}
|
||||
|
||||
/** Адрес WebSocket-канала организации: та же схема адресов, что у REST, и та
|
||||
* же сессия. Без активной организации канала нет — вернётся null. */
|
||||
export function resolveWebSocketUrl(path: string): string | null {
|
||||
if (typeof window === "undefined" || !activeOrganizationPublicId) return null;
|
||||
const base = configuredApiBase
|
||||
? new URL(configuredApiBase, window.location.origin)
|
||||
: new URL(window.location.origin);
|
||||
const protocol = base.protocol === "https:" ? "wss:" : "ws:";
|
||||
return `${protocol}//${base.host}/ws/organizations/${activeOrganizationPublicId}${path}`;
|
||||
}
|
||||
|
||||
function getCookie(name: string): string {
|
||||
const cookie = document.cookie
|
||||
.split("; ")
|
||||
|
||||
@@ -42,6 +42,7 @@ export function scopeLabel(scope: DialogScope): string {
|
||||
}
|
||||
import type { ConversationListItem, ListSort, ListTab } from "./types";
|
||||
import { useConversationCall } from "./useConversationCall";
|
||||
import { useConversationEvents } from "./useConversationEvents";
|
||||
import { useDebounced } from "../../shared/useDebounced";
|
||||
import { useConversationHistory } from "./useConversationHistory";
|
||||
import { useConversationList } from "./useConversationList";
|
||||
@@ -94,8 +95,19 @@ export function ConversationWorkspace({ isOwner = false, viewerId = null, listTi
|
||||
...(settledSearch ? { q: settledSearch } : {}),
|
||||
sort,
|
||||
}), [listTab, scope, settledSearch, sort]);
|
||||
const list = useConversationList(query);
|
||||
const history = useConversationHistory(selectedId);
|
||||
// Оповещения ведут обновление, опрос остаётся страховкой: при обрыве сокета
|
||||
// всё возвращается к прежним интервалам само.
|
||||
const events = useConversationEvents({
|
||||
conversationId: selectedId,
|
||||
onInboxChanged: () => void list.refresh(),
|
||||
onConversationChanged: (changedId) => {
|
||||
if (changedId !== selectedIdRef.current) return;
|
||||
void history.catchUp();
|
||||
void loadDetail(changedId);
|
||||
},
|
||||
});
|
||||
const list = useConversationList(query, { live: events.connected });
|
||||
const history = useConversationHistory(selectedId, { live: events.connected });
|
||||
|
||||
const loadDetail = useCallback(async (id: number) => {
|
||||
try {
|
||||
@@ -128,10 +140,11 @@ export function ConversationWorkspace({ isOwner = false, viewerId = null, listTi
|
||||
setActionError("");
|
||||
void loadDetail(selectedId);
|
||||
// Карточка диалога (статус, ответственный, метки) обновляется отдельно от
|
||||
// ленты: сообщений она больше не несёт.
|
||||
const timer = setInterval(() => loadDetail(selectedId), 3000);
|
||||
// ленты: сообщений она больше не несёт. С живыми оповещениями опрос — тоже
|
||||
// страховка.
|
||||
const timer = setInterval(() => loadDetail(selectedId), events.connected ? 30000 : 3000);
|
||||
return () => clearInterval(timer);
|
||||
}, [selectedId, loadDetail]);
|
||||
}, [events.connected, selectedId, loadDetail]);
|
||||
|
||||
const onConversationChanged = useCallback(() => {
|
||||
if (selectedId != null) void loadDetail(selectedId);
|
||||
|
||||
@@ -0,0 +1,91 @@
|
||||
import { useEffect, useRef, useState } from "react";
|
||||
|
||||
import { resolveWebSocketUrl } from "../../api/client";
|
||||
|
||||
// Оповещения о диалогах: сервер сообщает, что изменилось, клиент забирает
|
||||
// данные обычным запросом. Поллинг остаётся запасным путём и замедляется, пока
|
||||
// сокет жив, — при обрыве всё работает ровно как раньше.
|
||||
|
||||
const RECONNECT_MIN_MS = 1000;
|
||||
const RECONNECT_MAX_MS = 30000;
|
||||
|
||||
export type ConversationEvents = {
|
||||
/** Сокет открыт: поллинг можно замедлить. */
|
||||
connected: boolean;
|
||||
};
|
||||
|
||||
export function useConversationEvents({
|
||||
conversationId,
|
||||
onInboxChanged,
|
||||
onConversationChanged,
|
||||
}: {
|
||||
conversationId: number | null;
|
||||
onInboxChanged: () => void;
|
||||
onConversationChanged: (conversationId: number) => void;
|
||||
}): ConversationEvents {
|
||||
const [connected, setConnected] = useState(false);
|
||||
// Обработчики пересоздаются на каждый рендер — держим их в ref, чтобы сокет
|
||||
// не переоткрывался вместе с ними.
|
||||
const inboxRef = useRef(onInboxChanged);
|
||||
inboxRef.current = onInboxChanged;
|
||||
const conversationRef = useRef(onConversationChanged);
|
||||
conversationRef.current = onConversationChanged;
|
||||
const socketRef = useRef<WebSocket | null>(null);
|
||||
const watchedRef = useRef<number | null>(null);
|
||||
|
||||
useEffect(() => {
|
||||
const url = resolveWebSocketUrl("/conversations/");
|
||||
if (!url) return;
|
||||
let closed = false;
|
||||
let retry = RECONNECT_MIN_MS;
|
||||
let reconnectTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
|
||||
function open() {
|
||||
if (closed) return;
|
||||
const socket = new WebSocket(url as string);
|
||||
socketRef.current = socket;
|
||||
socket.onopen = () => {
|
||||
retry = RECONNECT_MIN_MS;
|
||||
setConnected(true);
|
||||
// После обрыва подписка теряется вместе с сокетом — восстанавливаем её.
|
||||
if (watchedRef.current != null) {
|
||||
socket.send(JSON.stringify({ type: "watch", conversationId: watchedRef.current }));
|
||||
}
|
||||
};
|
||||
socket.onmessage = (event) => {
|
||||
const payload = JSON.parse(String(event.data)) as { type?: string; conversationId?: number };
|
||||
if (payload.type === "inbox.changed") inboxRef.current();
|
||||
if (payload.type === "conversation.changed" && typeof payload.conversationId === "number") {
|
||||
conversationRef.current(payload.conversationId);
|
||||
}
|
||||
};
|
||||
socket.onclose = () => {
|
||||
setConnected(false);
|
||||
socketRef.current = null;
|
||||
if (closed) return;
|
||||
// Отступ растёт до полуминуты: сервер мог уйти на перезапуск.
|
||||
reconnectTimer = setTimeout(open, retry);
|
||||
retry = Math.min(retry * 2, RECONNECT_MAX_MS);
|
||||
};
|
||||
socket.onerror = () => socket.close();
|
||||
}
|
||||
|
||||
open();
|
||||
return () => {
|
||||
closed = true;
|
||||
if (reconnectTimer) clearTimeout(reconnectTimer);
|
||||
socketRef.current?.close();
|
||||
socketRef.current = null;
|
||||
};
|
||||
}, []);
|
||||
|
||||
useEffect(() => {
|
||||
watchedRef.current = conversationId;
|
||||
const socket = socketRef.current;
|
||||
if (socket && socket.readyState === WebSocket.OPEN && conversationId != null) {
|
||||
socket.send(JSON.stringify({ type: "watch", conversationId }));
|
||||
}
|
||||
}, [conversationId, connected]);
|
||||
|
||||
return { connected };
|
||||
}
|
||||
@@ -6,7 +6,9 @@ import { fetchMessages, type ApiMessage } from "./model";
|
||||
// прокрутке, вниз — дельтой обновления. Целиком лента не запрашивается никогда:
|
||||
// в диалоге может быть сколько угодно сообщений.
|
||||
const HISTORY_WINDOW = 50;
|
||||
// Пока оповещения живы, опрос — только страховка на случай потерянного события.
|
||||
const DELTA_INTERVAL_MS = 3000;
|
||||
const DELTA_IDLE_INTERVAL_MS = 30000;
|
||||
|
||||
export type ConversationHistory = {
|
||||
messages: ApiMessage[];
|
||||
@@ -26,7 +28,10 @@ function mergeNewer(current: ApiMessage[], incoming: ApiMessage[]): ApiMessage[]
|
||||
return fresh.length ? [...current, ...fresh] : current;
|
||||
}
|
||||
|
||||
export function useConversationHistory(conversationId: number | null): ConversationHistory {
|
||||
export function useConversationHistory(
|
||||
conversationId: number | null,
|
||||
{ live = false }: { live?: boolean } = {},
|
||||
): ConversationHistory {
|
||||
const [messages, setMessages] = useState<ApiMessage[]>([]);
|
||||
const [loaded, setLoaded] = useState(false);
|
||||
const [hasOlder, setHasOlder] = useState(false);
|
||||
@@ -88,9 +93,12 @@ export function useConversationHistory(conversationId: number | null): Conversat
|
||||
|
||||
useEffect(() => {
|
||||
if (conversationId == null || !loaded) return;
|
||||
const timer = setInterval(() => void catchUp(), DELTA_INTERVAL_MS);
|
||||
const timer = setInterval(
|
||||
() => void catchUp(),
|
||||
live ? DELTA_IDLE_INTERVAL_MS : DELTA_INTERVAL_MS,
|
||||
);
|
||||
return () => clearInterval(timer);
|
||||
}, [catchUp, conversationId, loaded]);
|
||||
}, [catchUp, conversationId, live, loaded]);
|
||||
|
||||
const loadOlder = useCallback(() => {
|
||||
const id = conversationRef.current;
|
||||
|
||||
@@ -11,6 +11,8 @@ import {
|
||||
// активность) и вклеивает её в уже загруженное, не сбрасывая прокрутку.
|
||||
const LIST_WINDOW = 30;
|
||||
const REFRESH_INTERVAL_MS = 4000;
|
||||
// С живыми оповещениями опрос остаётся только страховкой.
|
||||
const REFRESH_IDLE_INTERVAL_MS = 30000;
|
||||
|
||||
export type ConversationListState = {
|
||||
conversations: ApiConversation[];
|
||||
@@ -33,7 +35,10 @@ export function mergeHead(
|
||||
return [...head, ...tail.filter((conversation) => !fresh.has(conversation.id))];
|
||||
}
|
||||
|
||||
export function useConversationList(query: ConversationListQuery): ConversationListState {
|
||||
export function useConversationList(
|
||||
query: ConversationListQuery,
|
||||
{ live = false }: { live?: boolean } = {},
|
||||
): ConversationListState {
|
||||
const [conversations, setConversations] = useState<ApiConversation[]>([]);
|
||||
const [total, setTotal] = useState(0);
|
||||
const [loaded, setLoaded] = useState(false);
|
||||
@@ -86,9 +91,12 @@ export function useConversationList(query: ConversationListQuery): ConversationL
|
||||
}, [stableQuery]);
|
||||
|
||||
useEffect(() => {
|
||||
const timer = setInterval(() => void refresh(), REFRESH_INTERVAL_MS);
|
||||
const timer = setInterval(
|
||||
() => void refresh(),
|
||||
live ? REFRESH_IDLE_INTERVAL_MS : REFRESH_INTERVAL_MS,
|
||||
);
|
||||
return () => clearInterval(timer);
|
||||
}, [refresh]);
|
||||
}, [live, refresh]);
|
||||
|
||||
const loadMore = useCallback(() => {
|
||||
if (cursor == null || loadingMoreRef.current) return;
|
||||
|
||||
Reference in new issue
Block a user