Надёжная real-time доставка сообщений в чатах: SSE buffer, Last-Event-ID, polling fallback
This commit is contained in:
@@ -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', () => {
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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({
|
||||
|
||||
@@ -23,9 +23,19 @@ export interface Event {
|
||||
taskId?: number;
|
||||
}
|
||||
|
||||
export interface BufferedEvent extends Event {
|
||||
id: string;
|
||||
timestamp: number;
|
||||
}
|
||||
|
||||
export class EventBus {
|
||||
private connections: Map<string, SSEConnection> = new Map();
|
||||
private heartbeatInterval = 30000; // 30 seconds
|
||||
private eventIdCounter = 0;
|
||||
// Буферы последних событий для восстановления после reconnect
|
||||
private userBuffers = new Map<number, BufferedEvent[]>();
|
||||
private orgBuffers = new Map<number, BufferedEvent[]>();
|
||||
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<number, BufferedEvent[]>, 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<string, BufferedEvent>();
|
||||
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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user