From 7b830c658d3bc5522656a42ba3364b64072531e1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=98=D0=BB=D1=8C=D1=8F=D1=81=20=D0=A1=D1=83=D0=BB=D1=82?= =?UTF-8?q?=D0=B0=D0=BD=D0=BE=D0=B2?= Date: Wed, 8 Jul 2026 21:15:40 +0300 Subject: [PATCH] =?UTF-8?q?=D0=9D=D0=B0=D0=B4=D1=91=D0=B6=D0=BD=D0=B0?= =?UTF-8?q?=D1=8F=20real-time=20=D0=B4=D0=BE=D1=81=D1=82=D0=B0=D0=B2=D0=BA?= =?UTF-8?q?=D0=B0=20=D1=81=D0=BE=D0=BE=D0=B1=D1=89=D0=B5=D0=BD=D0=B8=D0=B9?= =?UTF-8?q?=20=D0=B2=20=D1=87=D0=B0=D1=82=D0=B0=D1=85:=20SSE=20buffer,=20L?= =?UTF-8?q?ast-Event-ID,=20polling=20fallback?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- client/src/components/TaskChat.tsx | 27 ++++++ client/src/hooks/useEvents.ts | 39 +++++--- client/src/lib/sseManager.ts | 22 ++++- client/src/pages/Chat.tsx | 27 ++++++ server/routes/auth.sse.routes.ts | 11 +++ server/routes/chat.messages.routes.ts | 34 ++++--- server/routes/messenger.messages.routes.ts | 6 +- server/routes/shared.ts | 108 ++++++++++++++++++--- server/storage.ts | 2 +- server/storage/tasks.storage.ts | 7 +- 10 files changed, 232 insertions(+), 51 deletions(-) diff --git a/client/src/components/TaskChat.tsx b/client/src/components/TaskChat.tsx index 633c3b8..c5c345c 100644 --- a/client/src/components/TaskChat.tsx +++ b/client/src/components/TaskChat.tsx @@ -270,6 +270,33 @@ const TaskChat = ({ taskId }: TaskChatProps) => { }, }); + // Lightweight polling fallback: даже если SSE молчит или событие потеряно, + // раз в 10 секунд запрашиваем только новые сообщения после последнего известного id. + useEffect(() => { + if (!taskId || !user) return; + const interval = setInterval(async () => { + if (document.hidden) return; + const currentData = queryClient.getQueryData(['/api/tasks', taskId, 'messages']) as any; + const currentMessages: TaskMessage[] = currentData?.messages || []; + const lastServerMessage = currentMessages.filter((m) => m.id > 0).pop(); + if (!lastServerMessage) return; + try { + const res = await apiRequest('GET', `/api/tasks/${taskId}/messages?afterId=${lastServerMessage.id}`); + const data = await res.json(); + if (!data.success || !Array.isArray(data.messages) || data.messages.length === 0) return; + queryClient.setQueryData(['/api/tasks', taskId, 'messages'], (old: any) => { + if (!old) return { messages: data.messages }; + const existingIds = new Set(old.messages.map((m: TaskMessage) => m.id)); + const merged = [...old.messages, ...data.messages.filter((m: TaskMessage) => !existingIds.has(m.id))]; + return { ...old, messages: merged }; + }); + } catch { + // ignore polling errors + } + }, 10000); + return () => clearInterval(interval); + }, [taskId, user]); + // Загрузка сообщений с умным polling - С reconciliation для сохранения optimistic updates const { data: messagesData, isLoading } = useQuery({ queryKey: ['/api/tasks', taskId, 'messages'], diff --git a/client/src/hooks/useEvents.ts b/client/src/hooks/useEvents.ts index 734b2f5..036f4f4 100644 --- a/client/src/hooks/useEvents.ts +++ b/client/src/hooks/useEvents.ts @@ -1,4 +1,4 @@ -import { useEffect, useState } from 'react'; +import { useEffect, useState, useRef } from 'react'; import { useAuth } from './useAuth'; import { sseManager } from '@/lib/sseManager'; @@ -21,6 +21,12 @@ export const useEvents = (options: UseEventsOptions = {}) => { const { user } = useAuth(); const [isConnected, setIsConnected] = useState(false); const [connectionId, setConnectionId] = useState(null); + // Храним актуальные callbacks в ref, чтобы слушатели SSE всегда вызывали + // последние версии обработчиков, но не переподписывались на каждый рендер. + const optionsRef = useRef(options); + useEffect(() => { + optionsRef.current = options; + }); // Привязываем singleton к текущему пользователю useEffect(() => { @@ -35,22 +41,25 @@ export const useEvents = (options: UseEventsOptions = {}) => { }); }, []); - // Подписываемся на события — каждый useEvents добавляет свои callbacks к одному EventSource - // Fixed: empty dependency array so we subscribe once on mount instead of - // unsubscribing/resubscribing on every render (which dropped SSE events). + // Подписываемся на события — каждый useEvents добавляет свои callbacks к одному EventSource. + // Регистрируем обёртки один раз; обёртки читают актуальные callbacks из ref. useEffect(() => { const unsubs: Array<() => void> = []; - if (options.onMessageRead) unsubs.push(sseManager.on('message_read', options.onMessageRead)); - if (options.onMessageCreated) unsubs.push(sseManager.on('message_created', options.onMessageCreated)); - if (options.onTaskUpdated) unsubs.push(sseManager.on('task_updated', options.onTaskUpdated)); - if (options.onNotification) unsubs.push(sseManager.on('notification', options.onNotification)); - if (options.onConvMessageCreated) unsubs.push(sseManager.on('conv_message_created', options.onConvMessageCreated)); - if (options.onConvMessageUpdated) unsubs.push(sseManager.on('conv_message_updated', options.onConvMessageUpdated)); - if (options.onConvMessageDeleted) unsubs.push(sseManager.on('conv_message_deleted', options.onConvMessageDeleted)); - if (options.onConvRead) unsubs.push(sseManager.on('conv_read', options.onConvRead)); - if (options.onConvCreated) unsubs.push(sseManager.on('conv_created', options.onConvCreated)); - if (options.onReactionUpdated) unsubs.push(sseManager.on('reaction_updated', options.onReactionUpdated)); - if (options.onPollUpdated) unsubs.push(sseManager.on('poll_updated', options.onPollUpdated)); + const wrap = (name: K) => (data: any) => { + const cb = optionsRef.current[name] as any; + if (cb) cb(data); + }; + if (options.onMessageRead) unsubs.push(sseManager.on('message_read', wrap('onMessageRead'))); + if (options.onMessageCreated) unsubs.push(sseManager.on('message_created', wrap('onMessageCreated'))); + if (options.onTaskUpdated) unsubs.push(sseManager.on('task_updated', wrap('onTaskUpdated'))); + if (options.onNotification) unsubs.push(sseManager.on('notification', wrap('onNotification'))); + if (options.onConvMessageCreated) unsubs.push(sseManager.on('conv_message_created', wrap('onConvMessageCreated'))); + if (options.onConvMessageUpdated) unsubs.push(sseManager.on('conv_message_updated', wrap('onConvMessageUpdated'))); + if (options.onConvMessageDeleted) unsubs.push(sseManager.on('conv_message_deleted', wrap('onConvMessageDeleted'))); + if (options.onConvRead) unsubs.push(sseManager.on('conv_read', wrap('onConvRead'))); + if (options.onConvCreated) unsubs.push(sseManager.on('conv_created', wrap('onConvCreated'))); + if (options.onReactionUpdated) unsubs.push(sseManager.on('reaction_updated', wrap('onReactionUpdated'))); + if (options.onPollUpdated) unsubs.push(sseManager.on('poll_updated', wrap('onPollUpdated'))); return () => { for (const u of unsubs) u(); }; }, []); diff --git a/client/src/lib/sseManager.ts b/client/src/lib/sseManager.ts index 8f1d5cb..4b82e65 100644 --- a/client/src/lib/sseManager.ts +++ b/client/src/lib/sseManager.ts @@ -12,6 +12,7 @@ let connectedListeners = new Set<(connected: boolean, connId: string | null) => let lastConnState = false; let lastConnId: string | null = null; let activeUserId: number | null = null; +let lastEventId: string | null = null; const KNOWN_EVENTS = [ 'message_read', @@ -68,7 +69,10 @@ async function connect() { } const { token } = await tokenRes.json(); - const es = new EventSource(`/api/events?token=${encodeURIComponent(token)}`); + // Передаём lastEventId явно для ручного reconnect (EventSource делает это + // автоматически только при собственном reconnect'е). + const lastEventIdParam = lastEventId ? `&lastEventId=${encodeURIComponent(lastEventId)}` : ''; + const es = new EventSource(`/api/events?token=${encodeURIComponent(token)}${lastEventIdParam}`); eventSource = es; isConnecting = false; @@ -80,6 +84,9 @@ async function connect() { es.onmessage = (event) => { try { + if (event.lastEventId) { + lastEventId = event.lastEventId; + } const data = JSON.parse(event.data); if (data.type === 'connected') { lastConnId = data.userId; @@ -91,10 +98,19 @@ async function connect() { for (const evName of KNOWN_EVENTS) { es.addEventListener(evName, (e: MessageEvent) => { try { + if (e.lastEventId) { + lastEventId = e.lastEventId; + } const data = JSON.parse(e.data); - console.log(`[SSE] Received event: ${evName}`, data); + if (process.env.NODE_ENV !== 'production') { + console.log(`[SSE] Received event: ${evName}`, data); + } emit(evName, data); - } catch {} + } catch (err) { + if (process.env.NODE_ENV !== 'production') { + console.error(`[SSE] Failed to parse event ${evName}:`, err); + } + } }); } diff --git a/client/src/pages/Chat.tsx b/client/src/pages/Chat.tsx index c144c1e..66bb314 100644 --- a/client/src/pages/Chat.tsx +++ b/client/src/pages/Chat.tsx @@ -1415,6 +1415,33 @@ export default function ChatPage() { }, []), }); + // Lightweight polling fallback: если SSE не доставил событие, раз в 10 секунд + // запрашиваем только новые сообщения после последнего известного id. + useEffect(() => { + if (!activeConvId || !user) return; + const interval = setInterval(async () => { + if (document.hidden) return; + const currentData = queryClient.getQueryData(['/api/messenger/conversations', activeConvId, 'messages']) as any; + const currentMessages: ConvMessage[] = currentData?.messages || []; + const lastServerMessage = currentMessages.filter((m) => m.id > 0).pop(); + if (!lastServerMessage) return; + try { + const res = await apiRequest('GET', `/api/messenger/conversations/${activeConvId}/messages?afterId=${lastServerMessage.id}&limit=50`); + const data = await res.json(); + if (!data.success || !Array.isArray(data.messages) || data.messages.length === 0) return; + queryClient.setQueryData(['/api/messenger/conversations', activeConvId, 'messages'], (old: any) => { + if (!old) return { success: true, messages: data.messages }; + const existingIds = new Set(old.messages.map((m: ConvMessage) => m.id)); + const merged = [...old.messages, ...data.messages.filter((m: ConvMessage) => !existingIds.has(m.id))]; + return { ...old, messages: merged }; + }); + } catch { + // ignore polling errors + } + }, 10000); + return () => clearInterval(interval); + }, [activeConvId, user]); + // ── Auto-scroll ─────────────────────────────────────────────────────────────── // Scroll to bottom whenever the active conversation changes (open from cache, new chat, etc.) diff --git a/server/routes/auth.sse.routes.ts b/server/routes/auth.sse.routes.ts index 947b3da..d45519f 100644 --- a/server/routes/auth.sse.routes.ts +++ b/server/routes/auth.sse.routes.ts @@ -78,6 +78,17 @@ export function registerSseRoutes(router: Router): void { }; eventBus.addConnection(connection); + + // Восстанавливаем события, пропущенные с момента lastEventId, до отправки connected. + // Это гарантирует, что клиент не потеряет сообщения во время reconnect. + // Поддерживаем как нативный заголовок EventSource, так и query-param для ручного reconnect. + const lastEventId = + (req.headers['last-event-id'] as string | undefined) || + (req.query.lastEventId as string | undefined); + if (lastEventId) { + eventBus.replayBufferedEvents(connection, lastEventId); + } + res.write(`data: {"type":"connected","userId":${user.id},"organizationId":${user.organizationId}}\n\n`); req.on('close', () => { diff --git a/server/routes/chat.messages.routes.ts b/server/routes/chat.messages.routes.ts index 2ecf902..ee88e71 100644 --- a/server/routes/chat.messages.routes.ts +++ b/server/routes/chat.messages.routes.ts @@ -41,7 +41,8 @@ export function registerChatMessageRoutes(router: Router): void { } } - const messages = await storage.getTaskMessages(taskId, req.organizationId!); + const afterId = req.query.afterId ? parseInt(req.query.afterId as string) : undefined; + const messages = await storage.getTaskMessages(taskId, req.organizationId!, afterId); res.json({ success: true, messages }); } catch (error) { console.error('Get task messages error:', error); @@ -178,8 +179,9 @@ export function registerChatMessageRoutes(router: Router): void { if (bot.type === 'ai_assistant') { if (!checkBotAccess(bot, req.user!)) { + let denialMessage; try { - await storage.createTaskMessage({ + denialMessage = await storage.createTaskMessage({ taskId, formId: task.formId, authorId: null, @@ -193,12 +195,20 @@ export function registerChatMessageRoutes(router: Router): void { } catch (err) { console.error(`[AI-Bot] Failed to post denial message for bot ${bot.id}:`, err); } - eventBus.publishEvent({ - type: 'message_created', - data: { taskId }, - organizationId: req.organizationId, - taskId, - }); + if (denialMessage) { + eventBus.publishEvent({ + type: 'message_created', + data: { + taskId, + message: { + ...denialMessage, + bot: bot ? { id: bot.id, name: bot.name, avatarUrl: bot.avatarUrl ?? null } : null, + }, + }, + organizationId: req.organizationId, + taskId, + }); + } continue; } @@ -209,14 +219,8 @@ export function registerChatMessageRoutes(router: Router): void { attachments: createdMessage.attachments || undefined, user: req.user!, organizationId: req.organizationId!, - }).then(() => { - eventBus.publishEvent({ - type: 'message_created', - data: { taskId }, - organizationId: req.organizationId, - taskId, - }); }).catch(err => console.error(`[AI-Bot] Error for bot ${bot.name}:`, err)); + // postBotReply внутри handleAiBotMention сам публикует message_created с полным сообщением continue; } diff --git a/server/routes/messenger.messages.routes.ts b/server/routes/messenger.messages.routes.ts index 6855195..a8350b8 100644 --- a/server/routes/messenger.messages.routes.ts +++ b/server/routes/messenger.messages.routes.ts @@ -2,7 +2,7 @@ import { Router } from "express"; import { z } from "zod"; import { db, openTenantCtx } from "../db"; import { conversations, conversationMembers, conversationMessages, users, bots, type Bot } from "@shared/schema"; -import { eq, and, lt, desc, sql, inArray } from "drizzle-orm"; +import { eq, and, lt, gt, desc, sql, inArray } from "drizzle-orm"; import { authenticateToken, type AuthenticatedRequest } from "../middleware/auth.middleware"; import { tenantIsolation } from "../middleware/tenant.middleware"; import { eventBus } from "./shared"; @@ -28,6 +28,7 @@ export function registerMessengerMessageRoutes(router: Router): void { const limit = Math.min(parseInt((req.query.limit as string) ?? "50"), 100); const beforeId = req.query.beforeId ? parseInt(req.query.beforeId as string) : null; + const afterId = req.query.afterId ? parseInt(req.query.afterId as string) : null; const conditions = [ eq(conversationMessages.conversationId, convId), @@ -36,6 +37,9 @@ export function registerMessengerMessageRoutes(router: Router): void { if (beforeId !== null && !isNaN(beforeId)) { conditions.push(lt(conversationMessages.id, beforeId)); } + if (afterId !== null && !isNaN(afterId)) { + conditions.push(gt(conversationMessages.id, afterId)); + } const msgs = await db .select({ diff --git a/server/routes/shared.ts b/server/routes/shared.ts index 1c61c28..f2323d3 100644 --- a/server/routes/shared.ts +++ b/server/routes/shared.ts @@ -23,9 +23,19 @@ export interface Event { taskId?: number; } +export interface BufferedEvent extends Event { + id: string; + timestamp: number; +} + export class EventBus { private connections: Map = new Map(); private heartbeatInterval = 30000; // 30 seconds + private eventIdCounter = 0; + // Буферы последних событий для восстановления после reconnect + private userBuffers = new Map(); + private orgBuffers = new Map(); + private readonly BUFFER_SIZE = 500; addConnection(connection: SSEConnection) { this.connections.set(connection.id, connection); @@ -54,6 +64,41 @@ export class EventBus { } } + private generateEventId(): string { + return `${Date.now()}-${++this.eventIdCounter}`; + } + + private pushToBuffer(map: Map, key: number, event: BufferedEvent) { + let buf = map.get(key); + if (!buf) { + buf = []; + map.set(key, buf); + } + buf.push(event); + if (buf.length > this.BUFFER_SIZE) { + buf.shift(); + } + } + + private writeEventToConnection(connection: SSEConnection, event: BufferedEvent): boolean { + try { + const eventData = JSON.stringify(event.data); + connection.res.write(`id: ${event.id}\n`); + connection.res.write(`event: ${event.type}\n`); + connection.res.write(`data: ${eventData}\n\n`); + // Force immediate delivery for SSE streams (Express/Node may buffer writes) + const resAny = connection.res as any; + if (typeof resAny.flush === 'function') { + resAny.flush(); + } + return true; + } catch (error) { + console.error(`Failed to send event to connection ${connection.id}:`, error); + this.removeConnection(connection.id); + return false; + } + } + publishEvent(event: Event) { // Also trigger an offline sync via silent web push when a task changes. // This wakes up closed PWAs so they can pull the latest data via delta sync. @@ -63,6 +108,19 @@ export class EventBus { .catch(() => {}); } + const bufferedEvent: BufferedEvent = { + ...event, + id: this.generateEventId(), + timestamp: Date.now(), + }; + + // Буферизуем событие для восстановления после reconnect + if (event.userId) { + this.pushToBuffer(this.userBuffers, event.userId, bufferedEvent); + } else if (event.organizationId) { + this.pushToBuffer(this.orgBuffers, event.organizationId, bufferedEvent); + } + this.connections.forEach((connection, id) => { // Проверяем tenant isolation if (event.organizationId && connection.organizationId !== event.organizationId) { @@ -74,22 +132,46 @@ export class EventBus { return; } - try { - const eventData = JSON.stringify(event.data); - connection.res.write(`event: ${event.type}\n`); - connection.res.write(`data: ${eventData}\n\n`); - // Force immediate delivery for SSE streams (Express/Node may buffer writes) - const resAny = connection.res as any; - if (typeof resAny.flush === 'function') { - resAny.flush(); - } - } catch (error) { - console.error(`Failed to send event to connection ${id}:`, error); - this.removeConnection(id); - } + this.writeEventToConnection(connection, bufferedEvent); }); } + /** + * Восстанавливает события, пропущенные после reconnect. + * Отправляет пользователю его персональные события + organization-wide события. + */ + replayBufferedEvents(connection: SSEConnection, lastEventId?: string) { + const userBuf = this.userBuffers.get(connection.userId) ?? []; + const orgBuf = this.orgBuffers.get(connection.organizationId) ?? []; + + // Объединяем без дубликатов через Map и сортируем по id (monotonic строка вида timestamp-counter) + const seen = new Map(); + for (const e of userBuf) seen.set(e.id, e); + for (const e of orgBuf) seen.set(e.id, e); + const combined = Array.from(seen.values()).sort((a, b) => a.id.localeCompare(b.id)); + + let startIdx = 0; + if (lastEventId) { + const idx = combined.findIndex(e => e.id === lastEventId); + if (idx !== -1) { + startIdx = idx + 1; + } + // Если lastEventId не найден в буфере, значит он слишком старый — + // отправляем всё, что есть (fallback, клиент сам разберёт дубликаты). + } + + const eventsToReplay = combined.slice(startIdx); + if (eventsToReplay.length === 0) return; + + for (const event of eventsToReplay) { + // Tenant isolation + if (event.organizationId && event.organizationId !== connection.organizationId) continue; + // User-specific events only for target user + if (event.userId && event.userId !== connection.userId) continue; + this.writeEventToConnection(connection, event); + } + } + getActiveConnections() { return this.connections.size; } diff --git a/server/storage.ts b/server/storage.ts index 2cdbd39..ea24030 100644 --- a/server/storage.ts +++ b/server/storage.ts @@ -209,7 +209,7 @@ export interface IStorage { createFieldHistory(insertHistory: InsertFieldHistory): Promise; // Task Messages (Chat) - getTaskMessages(taskId: number, organizationId: number): Promise; + getTaskMessages(taskId: number, organizationId: number, afterId?: number): Promise; getTaskMessage(messageId: number, organizationId: number): Promise; createTaskMessage(insertMessage: any, organizationId: number): Promise; updateTaskMessage(id: number, organizationId: number, updates: Partial): Promise; diff --git a/server/storage/tasks.storage.ts b/server/storage/tasks.storage.ts index 5ceeb5d..fed418b 100644 --- a/server/storage/tasks.storage.ts +++ b/server/storage/tasks.storage.ts @@ -2,7 +2,7 @@ import { users, forms, tasks, taskMessages, messageReads, bots, roles, type User import { taskRelations } from "@shared/schema"; import { taskAssignees, type TaskAssignee } from "@shared/schema"; import { db } from "../db"; -import { eq, and, desc, exists, sql } from "drizzle-orm"; +import { eq, and, desc, exists, sql, gt } from "drizzle-orm"; import { formatUserName } from "../utils/formatUserName"; import { TasksCoreStorage } from "./tasks-core.storage"; import { fieldConditionsMatch, resolveFieldDisplayText } from "../utils/field-conditions"; @@ -239,7 +239,7 @@ export class TasksStorage extends TasksCoreStorage { } // Task Messages (Chat) - async getTaskMessages(taskId: number, organizationId: number): Promise { + async getTaskMessages(taskId: number, organizationId: number, afterId?: number): Promise { const replyToMsg = db.select({ id: taskMessages.id, message: taskMessages.message, @@ -289,7 +289,8 @@ export class TasksStorage extends TasksCoreStorage { .leftJoin(replyToAuthor, eq(replyToMsg.authorId, replyToAuthor.id)) .where(and( eq(taskMessages.taskId, taskId), - eq(forms.organizationId, organizationId) + eq(forms.organizationId, organizationId), + ...(afterId ? [gt(taskMessages.id, afterId)] : []) )) .orderBy(taskMessages.createdAt);