From 03825dfbd5acf7e3a5c8ac6c028b5bc33a8992cf Mon Sep 17 00:00:00 2001 From: Andrey Date: Wed, 9 Sep 2026 02:33:28 +0300 Subject: [PATCH] =?UTF-8?q?:zap:=20perf(realtime):=20=D0=BE=D0=B1=D0=BD?= =?UTF-8?q?=D0=BE=D0=B2=D0=BB=D0=B5=D0=BD=D0=B8=D0=B5=20=D0=BF=D0=BE=20?= =?UTF-8?q?=D1=81=D0=BE=D0=B1=D1=8B=D1=82=D0=B8=D1=8E=20=D0=B2=D0=BC=D0=B5?= =?UTF-8?q?=D1=81=D1=82=D0=BE=20=D0=BE=D0=BF=D1=80=D0=BE=D1=81=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Клиент опрашивал сервер: список раз в четыре секунды, карточку и дельту истории — раз в три; на оператора выходило под сорок запросов в минуту, и новое сообщение всё равно появлялось с задержкой. - WebSocket-канал организации (channels уже стоял ради сигналинга звонков): аутентификация — сессией того же SPA, организация — в адресе, как в HTTP; - событие несёт только повод обновиться, данные клиент забирает обычным запросом: проверка видимости остаётся в одном месте, и канал не может в ней ошибиться. Событие инбокса не содержит идентификаторов — сотрудник видит не все диалоги организации; на события диалога подписка отдельная, и сервер проверяет видимость перед ней; - события шлют те же сигналы, что держат свежесть диалога: одно место на семь мест создания сообщений и на все изменения состояния; - опрос остался запасным путём и замедляется до тридцати секунд, пока канал жив: при обрыве всё возвращается к прежнему поведению само. Проверено сквозь шлюз: апгрейд проходит до консьюмера, анонимное соединение отклоняется кодом 403. Co-Authored-By: Claude Opus 5 --- .../chatballs/conversations/consumers.py | 116 ++++++++++++ .../chatballs/conversations/realtime.py | 59 ++++++ .../chatballs/conversations/routing.py | 14 ++ .../chatballs/conversations/signals.py | 37 +++- .../chatballs/conversations/test_realtime.py | 168 ++++++++++++++++++ apps/backend/chatballs_backend/asgi_app.py | 7 +- apps/internal-ui/src/api/client.ts | 11 ++ .../conversations/ConversationWorkspace.tsx | 23 ++- .../conversations/useConversationEvents.ts | 91 ++++++++++ .../conversations/useConversationHistory.ts | 14 +- .../conversations/useConversationList.ts | 14 +- 11 files changed, 536 insertions(+), 18 deletions(-) create mode 100644 apps/backend/chatballs/conversations/consumers.py create mode 100644 apps/backend/chatballs/conversations/realtime.py create mode 100644 apps/backend/chatballs/conversations/routing.py create mode 100644 apps/backend/chatballs/conversations/test_realtime.py create mode 100644 apps/internal-ui/src/features/conversations/useConversationEvents.ts diff --git a/apps/backend/chatballs/conversations/consumers.py b/apps/backend/chatballs/conversations/consumers.py new file mode 100644 index 0000000..f43b760 --- /dev/null +++ b/apps/backend/chatballs/conversations/consumers.py @@ -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) diff --git a/apps/backend/chatballs/conversations/realtime.py b/apps/backend/chatballs/conversations/realtime.py new file mode 100644 index 0000000..f0b6c2d --- /dev/null +++ b/apps/backend/chatballs/conversations/realtime.py @@ -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) diff --git a/apps/backend/chatballs/conversations/routing.py b/apps/backend/chatballs/conversations/routing.py new file mode 100644 index 0000000..02364bf --- /dev/null +++ b/apps/backend/chatballs/conversations/routing.py @@ -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//conversations/", + AuthMiddlewareStack(ConversationEventsConsumer.as_asgi()), + ), +] diff --git a/apps/backend/chatballs/conversations/signals.py b/apps/backend/chatballs/conversations/signals.py index 1034ba9..dbe4588 100644 --- a/apps/backend/chatballs/conversations/signals.py +++ b/apps/backend/chatballs/conversations/signals.py @@ -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) diff --git a/apps/backend/chatballs/conversations/test_realtime.py b/apps/backend/chatballs/conversations/test_realtime.py new file mode 100644 index 0000000..692d427 --- /dev/null +++ b/apps/backend/chatballs/conversations/test_realtime.py @@ -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() diff --git a/apps/backend/chatballs_backend/asgi_app.py b/apps/backend/chatballs_backend/asgi_app.py index b26f695..212400f 100644 --- a/apps/backend/chatballs_backend/asgi_app.py +++ b/apps/backend/chatballs_backend/asgi_app.py @@ -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), } ) diff --git a/apps/internal-ui/src/api/client.ts b/apps/internal-ui/src/api/client.ts index 27042b9..a6aed6f 100644 --- a/apps/internal-ui/src/api/client.ts +++ b/apps/internal-ui/src/api/client.ts @@ -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("; ") diff --git a/apps/internal-ui/src/features/conversations/ConversationWorkspace.tsx b/apps/internal-ui/src/features/conversations/ConversationWorkspace.tsx index aad2c8d..8301f49 100644 --- a/apps/internal-ui/src/features/conversations/ConversationWorkspace.tsx +++ b/apps/internal-ui/src/features/conversations/ConversationWorkspace.tsx @@ -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); diff --git a/apps/internal-ui/src/features/conversations/useConversationEvents.ts b/apps/internal-ui/src/features/conversations/useConversationEvents.ts new file mode 100644 index 0000000..3902f02 --- /dev/null +++ b/apps/internal-ui/src/features/conversations/useConversationEvents.ts @@ -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(null); + const watchedRef = useRef(null); + + useEffect(() => { + const url = resolveWebSocketUrl("/conversations/"); + if (!url) return; + let closed = false; + let retry = RECONNECT_MIN_MS; + let reconnectTimer: ReturnType | 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 }; +} diff --git a/apps/internal-ui/src/features/conversations/useConversationHistory.ts b/apps/internal-ui/src/features/conversations/useConversationHistory.ts index a2f4854..959ae2d 100644 --- a/apps/internal-ui/src/features/conversations/useConversationHistory.ts +++ b/apps/internal-ui/src/features/conversations/useConversationHistory.ts @@ -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([]); 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; diff --git a/apps/internal-ui/src/features/conversations/useConversationList.ts b/apps/internal-ui/src/features/conversations/useConversationList.ts index 6d0dcdd..28f4ac9 100644 --- a/apps/internal-ui/src/features/conversations/useConversationList.ts +++ b/apps/internal-ui/src/features/conversations/useConversationList.ts @@ -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([]); 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;